Skip to content

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

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

FrankChen021 wants to merge 2 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)

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.

3 participants