Skip to content

feat: support ingesting Kafka messages from selected partitions - #20474

Open
FrankChen021 wants to merge 7 commits into
apache:masterfrom
FrankChen021:kafka-partition-selection-upstream
Open

FrankChen021 wants to merge 7 commits into
apache:masterfrom
FrankChen021:kafka-partition-selection-upstream

Conversation

@FrankChen021

@FrankChen021 FrankChen021 commented Oct 2, 2026 •

Copy link
Copy Markdown
Member

Description

Adds an optional ioConfig.partitionIds to 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

  • partitionIds must be a nonempty set of nonnegative integers. It can't be combined with topicPattern or boundedStreamConfig. IDs may be JSON integers or numeric strings, and duplicates are normalized.
  • IDs are validated against the broker's 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.
  • Selected partitions are assigned to task groups by their rank in the sorted selection, so [0, 3, 6] with 3 tasks uses 3 groups. Without a selection, partition % taskCount is unchanged.
  • Changing or removing the selection is a spec update, which restarts the supervisor and its tasks. Offsets saved for excluded partitions are kept, so a partition resumes from them if it is re-added.
  • Lag, idle detection and resetToLatestAndBackfill look only at the selected partitions. A backfill spec never inherits the selection.
  • Partition-scoped resets (reset / resetOffsets) that name a different topic or any excluded partition are rejected. A full reset is still allowed and clears all saved offsets.
  • The web console ingestion form has a "Partition IDs" field which is by default collapsed, when expanded, users can fill partition ids
image

Release note

Kafka supervisors accept ioConfig.partitionIds to ingest only selected partitions of a topic, for example for independent diagnostic ingestion.


Key changed/added classes in this PR
  • KafkaSupervisorIOConfig
  • KafkaSupervisor
  • KafkaRecordSupplier
  • SeekableStreamSupervisor
  • SupervisorManager

This PR has:

  • been self-reviewed.
  • added documentation for new or modified features or behaviors.
  • a release note entry in the PR description.
  • added Javadocs for most classes and all non-trivial methods. Linked related entities via Javadoc links.
  • added comments explaining the "why" and the intent of the code wherever would not be obvious for an unfamiliar reader.
  • added unit tests or modified existing tests to cover new code paths, ensuring the threshold for code coverage is met.
  • added integration tests.
  • been tested in a test Druid cluster.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 Medium severity

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 partitionIds configuration 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.

Comment thread web-console/src/druid-models/ingestion-spec/ingestion-spec.tsx Outdated

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟢 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 FrankChen021 left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟢 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 FrankChen021 left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟢 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 amaechler left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Some comments from a Claude review to consider.

ioConfig.getServerPriorityToReplicas(),
boundedStreamConfig
boundedStreamConfig,
null

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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),

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 =>

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@FrankChen021
FrankChen021 requested a review from amaechler October 9, 2026 08:09

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟢 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)

@amaechler amaechler left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐻

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants