feat: support ingesting Kafka messages from selected partitions - #20474
FrankChen021 wants to merge 2 commits into
Conversation
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
The web console rejects valid numeric-string IDs and accepts out-of-range IDs that the server rejects.
Review effort: Balanced
Findings: 1
Open (1)
What changed in this PR
Adds optional Kafka partition selection across supervisor ingestion, sampling, lag monitoring, resets, backfills, documentation, and the web console.
Changes:
- Adds validated
partitionIdsconfiguration and partition-aware task assignment. - Filters discovery, lag, sampling, resets, and backfills to selected partitions.
- Adds comprehensive unit, integration, console, and serialization tests.
| File | Description |
|---|---|
web-console/src/druid-models/ingestion-spec/ingestion-spec.tsx |
Adds the partition IDs form field and validation. |
web-console/src/druid-models/ingestion-spec/ingestion-spec.spec.ts |
Tests console field conversion and validation. |
processing/src/main/java/org/apache/druid/jackson/StrictIntegerDeserializer.java |
Adds strict integer deserialization. |
processing/src/test/java/org/apache/druid/jackson/StrictIntegerDeserializerTest.java |
Tests accepted and rejected integer representations. |
indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java |
Adds reset validation and partition-selection extension points. |
indexing-service/src/main/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManager.java |
Restricts backfill offsets to current partitions. |
indexing-service/src/test/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManagerTest.java |
Tests partition-filtered backfill offsets. |
extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/KafkaRecordSupplier.java |
Filters and validates discovered Kafka partitions. |
extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/KafkaSamplerSpec.java |
Applies selection during sampling. |
extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/supervisor/KafkaSupervisor.java |
Implements grouping, reset validation, and offset filtering. |
extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/supervisor/KafkaSupervisorIOConfig.java |
Defines and validates partitionIds. |
extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/supervisor/KafkaIOConfigBuilder.java |
Adds builder support for partition selection. |
extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/supervisor/KafkaSupervisorSpec.java |
Excludes selection from bounded backfills. |
extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/KafkaRecordSupplierTest.java |
Tests discovery filtering and missing partitions. |
extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/KafkaSamplerSpecTest.java |
Tests selected-partition sampling. |
extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/supervisor/KafkaSupervisorTest.java |
Tests grouping, lag, compatibility, and resets. |
extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/supervisor/KafkaSupervisorIOConfigTest.java |
Tests configuration validation, serialization, and updates. |
extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/supervisor/KafkaSupervisorSpecTest.java |
Tests backfill selection removal. |
extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/EmbeddedKafkaSupervisorTest.java |
Verifies end-to-end selected-partition ingestion. |
docs/ingestion/kafka-ingestion.md |
Documents configuration and operational behavior. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
…ccepted range and string forms
FrankChen021
left a comment
There was a problem hiding this comment.
🟢 Approval recommended
Reviewed 20 of 20 changed files.
No actionable issues found in this review. The selected-partition behavior is carried through supervisor configuration, Kafka discovery, sampling, task grouping, lag and idle calculations, reset validation, and reset-to-latest backfills. The current head also restarts supervisors and tasks when the selection changes, preserves excluded offsets for later re-addition, and validates the console's integer and range semantics consistently with the server. I inspected the surrounding seekable-stream lifecycle and task-adoption paths for compatibility and stale-task behavior.
Validation: git diff --check 3e9365196f9893ff9b1cefcfbcb57d49e4080b4a d6c219a7c4bbee7b1dc851e6e28822ebded0e37f passed. No builds, test suites, dependency installs, or formatters were run per the scoped static-review instructions.
This is an automated review by Codex GPT-5.6-Luna(max)

Description
Adds an optional
ioConfig.partitionIdsto single-topic Kafka supervisors, so operators can ingest a chosen subset of a topic's partitions, for example into a separate diagnostic datasource, without re-ingesting the whole topic. Omitting it keeps the current all-partition behavior.Behavior
partitionIdsmust be a nonempty set of nonnegative integers. It can't be combined withtopicPatternorboundedStreamConfig. IDs may be JSON integers or numeric strings, and duplicates are normalized.partitionsFor(topic)during discovery (supervisor, lag monitoring, sampler), not at spec parsing. A missing ID raises a supervisor error and is retried, so a supervisor never starts with part of its selection. Partitions added to the topic later that are outside the selection are ignored.[0, 3, 6]with 3 tasks uses 3 groups. Without a selection,partition % taskCountis unchanged.resetToLatestAndBackfilllook only at the selected partitions. A backfill spec never inherits the selection.reset/resetOffsets) that name a different topic or any excluded partition are rejected. A full reset is still allowed and clears all saved offsets.Release note
Kafka supervisors accept
ioConfig.partitionIdsto ingest only selected partitions of a topic, for example for independent diagnostic ingestion.Key changed/added classes in this PR
KafkaSupervisorIOConfigKafkaSupervisorKafkaRecordSupplierSeekableStreamSupervisorSupervisorManagerThis PR has: