Repository navigation
feat: support ingesting Kafka messages from selected partitions - #20474
FrankChen021 wants to merge 7 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)
…rtition-selection-upstream # Conflicts: # extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/KafkaRecordSupplier.java # extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/KafkaSamplerSpec.java # extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/supervisor/KafkaSupervisor.java # extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/supervisor/KafkaSupervisorIOConfig.java
FrankChen021
left a comment
There was a problem hiding this comment.
🟢 Approval recommended
No actionable issues found. The selection is normalized and validated, applied consistently to Kafka partition discovery and sampling, mapped deterministically across task groups, and carried through offset and lag reporting, reset validation, and reset-to-latest backfill bounds. Selection changes restart the relevant readers; omitting the selection retains all-partition behavior. I rechecked the previous console validation concern: the current validator accepts integer strings and rejects empty, negative, fractional, and out-of-range IDs.
Coverage: reviewed all 20 files in the current PR diff against the merged master parent. I also checked the incremental merge update and the overlap files, while excluding upstream-only changes from the PR review scope.
Validation: git diff --check HEAD^2 HEAD passed. No tests or builds were run; this was a scoped static review.
This is an automated review by Codex GPT-5.6-Luna(max)
FrankChen021
left a comment
There was a problem hiding this comment.
🟢 Approval recommended
No actionable issues found in the current head. The incremental diff contains only upstream changes from the latest master merge; I rechecked the complete current PR diff and surrounding selection, task lifecycle, reset, backfill, sampling, and console-validation paths, with no regression from the merge.
Reviewed 20 of 20 changed files in the current PR diff.
Validation: git diff --check origin/master...HEAD passed. No tests or builds were run; this was a static follow-up review.
This is an automated review by Codex GPT-5.6-Luna(max)
amaechler
left a comment
There was a problem hiding this comment.
Some comments from a Claude review to consider.
| ioConfig.getServerPriorityToReplicas(), | ||
| boundedStreamConfig | ||
| boundedStreamConfig, | ||
| null |
There was a problem hiding this comment.
Backfilling a supervisor that has partitionIds set seems to break the backfill's lag refresh. createBackfillSpec passes null for partitionIds, so the bounded supervisor's updatePartitionLagFromStream asks getPartitionIds for the whole topic. It then seeks and reads positions for partitions that bounded mode never assigned.
I reproduced this with a KafkaSupervisorTest case: a bounded config covering partition 0 of a 3-partition topic, then updatePartitionLagFromStream() throws NoOffsetForPartition for partitions 1 and 2. In a real cluster updateCurrentAndLatestOffsets catches it, so I suspect the backfill supervisor logs a warning every cycle and never reports lag (not confirmed end to end).
Options: let the backfill carry the selection, which means relaxing the boundedStreamConfig exclusion for derived specs. Or have the Kafka lag path use the partitions the supervisor actually tracks instead of re-querying the topic.
There was a problem hiding this comment.
Thanks, reproduced with a test (testBoundedSupervisorReportsLagOnlyForBoundedPartitions): a bounded config covering partition 0 of a 3-partition topic made updatePartitionLagFromStream() throw NoOffsetForPartition for partitions 1 and 2. Fixed in f9143b9. Both Kafka lag methods now take their partitions from getPartitionsForLag(), which returns the partitions the supervisor tracks (from the bounded config) when it is bounded, and the discovered topic partitions otherwise. I did not let the backfill carry the selection, so the boundedStreamConfig exclusion is unchanged.
| name: 'partitionIds', | ||
| label: 'Partition IDs', | ||
| type: 'string-array', | ||
| defined: ioConfig => ioConfig.type === 'kafka' && Boolean(ioConfig.topic), |
There was a problem hiding this comment.
If a user fills Partition IDs, then clears topic and enters a topicPattern, the Partition IDs field disappears. issueWithIoConfig still returns "partitionIds requires a single topic", so the form is stuck unless they edit the JSON. Could the field stay visible while it holds a value, or could partitionIds be cleared when topic is unset?
There was a problem hiding this comment.
Fixed in d73570f. defined is now type === 'kafka' && (Boolean(topic) || partitionIds != null), so the field stays visible while it holds a value and can be cleared after switching to topicPattern. There is a test for it in ingestion-spec.spec.ts.
| defined: ioConfig => ioConfig.type === 'kafka' && Boolean(ioConfig.topic), | ||
| placeholder: 'Optional; comma-separated, e.g. 0, 2', | ||
| hideInMore: ioConfig => ioConfig.partitionIds == null, | ||
| valueAdjustment: partitionIds => |
There was a problem hiding this comment.
Could this valueAdjustment just go? The server and isValidKafkaPartitionId already accept integer strings. map(Number) also loosens what the form accepts: 1e3 becomes 1000 and 0x10 becomes 16, while the server rejects both as strings.
There was a problem hiding this comment.
Agreed that map(Number) was too loose. Fixed in d73570f, but I kept the adjustment instead of removing it, so the spec JSON stays numeric ([0, 2] rather than ["0","2"]). Only plain decimal digit strings are converted; anything else (1e3, 0x10, -1, abc) stays a string and is rejected by issueWithIoConfig, as on the server. Tests cover this.
| * Retains the signature used before {@code partitionIds} was introduced (defaults it to null), so that callers | ||
| * compiled against it keep working. | ||
| */ | ||
| public KafkaSupervisorIOConfig( |
There was a problem hiding this comment.
Nit: this constructor has no callers outside tests, and the same goes for the two new null-passing overloads in KafkaRecordSupplier. Adding the parameter at the existing call sites would drop roughly 70 lines. It would also stop a future caller from silently losing partitionIds.
There was a problem hiding this comment.
Agreed. Removed in d73570f: the extra KafkaSupervisorIOConfig constructor and the two null-passing KafkaRecordSupplier overloads. The call sites (KafkaIndexTask and the tests) now pass partitionIds explicitly.
|
|
||
| /** Keeps offsets for the configured partitions only. */ | ||
| @Override | ||
| protected Map<KafkaTopicPartition, Long> getHighestCurrentOffsets() |
There was a problem hiding this comment.
Optional simplification: KafkaSupervisor.getOffsetsFromMetadataStorage() already filters stored offsets by topic. Adding the selection filter there would let you drop the new getOffsetsFromMetadataStorageForCurrentPartitions hook in SeekableStreamSupervisor, its two call-site changes and this override. As far as I can tell, every caller only needs the selected partitions. Offsets for excluded partitions are still kept, because commits and resets go through the metadata coordinator rather than this getter.
There was a problem hiding this comment.
Agreed, and I checked the other readers of that getter (task creation, bounded completion, highest-offsets report, idle check, backfill start offsets): all only need the selected partitions, and resets do not go through it. Done in 8f3b6b3. KafkaSupervisor.getOffsetsFromMetadataStorage() now also drops excluded partitions, and getOffsetsFromMetadataStorageForCurrentPartitions, its base-class hook, and my changes to SupervisorManager are gone. The filter in getHighestCurrentOffsets stays because publishing tasks can still carry excluded partitions.
…rop unused constructors
…e and drop the extra hook
FrankChen021
left a comment
There was a problem hiding this comment.
🟢 Approval recommended
No actionable issues found in the current head. The update filters stored offsets through the common Kafka metadata path while retaining excluded offsets in storage, and bounds lag polling to partitions assigned by the bounded configuration. I rechecked discovery, sampling, grouping, task adoption and replacement, reset validation, backfill construction, configuration serialization, and console validation; no merge-blocking regression was found.
Reviewed 20 of 20 changed files in the full current PR diff, starting with the incremental diff. Also inspected the two incremental-only reversion files in SupervisorManager and its tests, and surrounding seekable-stream lifecycle code.
Validation: git diff --check 131989817271bc836e815db738cc2be33fba3538 HEAD passed. This was a static review; tests and builds were not run.
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: