test(gateway): Kafka + gateway -> FiloDB ingestion integration test - #2342
Open
siri-varma wants to merge 6 commits into
Open
siri-varma wants to merge 6 commits into
siri-varma wants to merge 6 commits into
Conversation
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.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What
Adds an automated integration test (
gateway/src/it, the sbtIntegrationTestconfig) that stands up the real production ingestion stack and verifies FiloDB ingests through it: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
partition == shard).FilodbClusterconsumes viaKafkaIngestionStreamFactory(theIngestionStreamSpecpattern), no Cassandra.memStore.numRowsIngested(ref)until it matches the produced count.auto.offset.reset=earliestis set so produce/consume ordering cannot cause a false negative (KafkaIngestionStreamusesassign()+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— enableIntegrationTeston the gateway module.project/Dependencies.scala— testcontainers + scalaTest initscope..github/workflows/scala.yml— newkafka-gateway-integrationjob (JDK 21) runningsbt "gateway/it:test".No existing tests or the query path are modified.