Skip to content

test(gateway): Kafka + gateway -> FiloDB ingestion integration test - #2342

Open
siri-varma wants to merge 6 commits into
developfrom
users/svegiraju/kafka-gateway-ingest
Open

siri-varma wants to merge 6 commits into
developfrom
users/svegiraju/kafka-gateway-ingest

Conversation

@siri-varma

Copy link
Copy Markdown
Contributor

What

Adds an automated integration test (gateway/src/it, the sbt IntegrationTest config) that stands up the real production ingestion stack and verifies FiloDB ingests through it:

gateway data generators  ->  Kafka (testcontainers)  ->  FiloDB (in-process
FilodbCluster via KafkaIngestionStreamFactory)  ->  memStore.numRowsIngested

Covers gauge and OTel-exponential-histogram data produced by the gateway's own generators (TestTimeseriesProducer). This is an ingestion test — it verifies the pipeline delivers data and the ingested row count matches what was produced. It is not a PromQL-comparison test.

How

  • Testcontainers starts a real Kafka broker; the topic is created with one partition per shard (the gateway publishes partition == shard).
  • An in-process FilodbCluster consumes via KafkaIngestionStreamFactory (the IngestionStreamSpec pattern), no Cassandra.
  • Verification polls memStore.numRowsIngested(ref) until it matches the produced count.
  • auto.offset.reset=earliest is set so produce/consume ordering cannot cause a false negative (KafkaIngestionStream uses assign() + seek(offset) only when a checkpoint exists; with none it falls back to the reset policy).

Changes

  • gateway/src/it/scala/filodb/gateway/KafkaGatewayIngestionSpec.scala — the spec.
  • project/FiloBuild.scala — enable IntegrationTest on the gateway module.
  • project/Dependencies.scala — testcontainers + scalaTest in it scope.
  • .github/workflows/scala.yml — new kafka-gateway-integration job (JDK 21) running sbt "gateway/it:test".

No existing tests or the query path are modified.

Siri Varma Vegiraju added 6 commits September 9, 2026 14:33
Add an sbt IntegrationTest (gateway/src/it) that stands up the real
production ingestion stack and verifies FiloDB ingests through it:

  gateway generators -> Kafka (testcontainers) -> FiloDB (in-process
  FilodbCluster via KafkaIngestionStreamFactory) -> numRowsIngested

Covers gauge and OTel-exponential-histogram data produced by the
gateway's own generators. Verified via memStore.numRowsIngested; uses
auto.offset.reset=earliest so produce/consume ordering is irrelevant.

- project/FiloBuild.scala: enable IntegrationTest on the gateway module
- project/Dependencies.scala: testcontainers + scalaTest in it scope
- .github/workflows/scala.yml: kafka-gateway-integration job (JDK 21)

The CSV/PromCompat harness is untouched.
… consumer

The first CI run failed with 0 rows ingested despite consumers correctly
assigned + reset to offset 0 (auto.offset.reset=earliest works). It is
unclear whether the gateway producer never flushed to the topic, or the
records were consumed but dropped on decode.

Add condition-based waits on both boundaries and surface both counts in the
failure clue so the next run is self-diagnosing:
  - container-messages on the Kafka topic (producer side), and
  - numRowsIngested (consumer/decode side).
Root cause of 0-rows-ingested (confirmed by the diagnostic: containers were
on the topic, but numRowsIngested stayed 0): the gateway encodes records
with the global Schemas object (filodb-defaults predefined-keys, 10 keys),
while application_test.conf overrides partition-schema.predefined-keys to []
(a HOCON list replaces rather than merges). With mismatched tables the shard
cannot decode the compacted tag keys and silently drops every record.

Override the test's predefined-keys to filodb-defaults' 10-key list so the
encode and decode tables are identical. Consumers were already correctly
positioned (auto.offset.reset=earliest, offset 0) and the producer was
writing containers to the topic; only decode was failing.
core/src/test/resources/logback-test.xml is on the gateway it classpath and
sends the filodb logger to a FILE (additivity=false), so ingestion diagnostics
never reached CI stdout. Force a uniquely-named logback config via
-Dlogback.configurationFile so it can't be shadowed, and DEBUG the ingest path
(schema resolution, container version/empty checks, row counts) to stdout.
… ingest

Root cause (confirmed via it-run DEBUG logs): SetupDataset defaults
overrideSchema=true, so the shard registered only the dataset's single schema
(promCounter). The gateway generators emit gauge (schemaId 2200) and
otel-exp-delta-histogram (35138) records, which the shard then dropped as
'unknown schema' — 0 rows ingested. Pass overrideSchema=false so the shard
registers all config schemas (same filodb-defaults definitions the gateway
encodes with), matching the real production setup.
…ctly

Adds a third it-test that validates histogram *correctness*, not just row
counts. It hand-builds a deterministic exponential-histogram dataset (known
bucket values), publishes it through the real gateway sharding + Kafka path,
reads the raw 'h' column back with a MultiSchemaPartitionsExec scan, and
asserts each (timestamp, histogram) matches. Histogram.equals is scheme+values,
so the assertion fails if the bucket scheme or any bucket count is corrupted
anywhere in encode -> Kafka -> decode -> store.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant