apache / apache/druid

Range partitioning for native parallel batch indexing

Open
#8,769 2 comments 0 reactions 0 assignees View on GitHub
Design Review Proposal
Dominant language
Java
Stars
14.1k
Forks
3.8k
Avg merge
2d 58m
Merged PRs (30d)
233

Description

### Motivation

Currently, native parallel batch indexing (#5543) only supports hashed partitions. Add support for range partitioning so that native parallel batch indexing has feature parity with hadoop batch indexing.

### Proposed changes

For native parallel batch indexing hash partitioning, `ParallelIndexSupervisorTask.runMultiPhaseParallel()` runs two phases:

1. `PartialSegmentGenerateTask`: Each worker takes a firehose split and creates partial segments by hashing rows from the split
2. `PartialSegmentMergeTask`: Each worker merges the set of partial segments for each partition

For range partitions, the `PartialSegmentMergeTask` can be reused, but partial segment generation needs a different algorithm:

1. Each worker takes a firehose split and generates an approximate histogram in memory that records the frequency of each key (timestamp adjusted by query granularity and dimension value tuple) in the split. In particular, `quantiles.ItemsSketch` from DataSketches can be used as the implementation of the approximate histogram.

- The most accurate `ItemsSketch` of `k=32768` (which also uses more space) has a normalized rank error of 0.009% and for extremely large streams (>1B items), retains no more than 600k items. Assuming each item is about 1KB (~500 characters in the dimension value), the histogram will use a max of about 600MB memory, so it would not need to be periodically flushed to disk. I’m not sure what serialized size this translates to, but it’s possible that it’s too large to add as part of the task report payload, in which case it’d need to be written to disk and then copied to the supervisor.
- `k` can be exposed as a configurable value in the ingestion spec, in case users wish to reduce the memory consumption
- **NOTE:** Using `ItemsSketch` will require moving DataSketches from extensions to core, which may block https://github.com/apache/incubator-druid/pull/8647

2. The supervisor gets all the partial `ItemsSketch`es generated earlier (either via the `TaskReport` payload or by copying the serialized `ItemSketch`es from disk) and then merges them by using `quantiles.ItemsUnion`. From the merged `ItemsSketch`, it uses `ItemsSketch.getQuantiles()` to creates partition ranges that satisfy the criteria specified in `SingleDimensionPartitionsSpec`.

- Because `ItemsSketch` yields estimated quantile values, enforcing `maxRowsPerSegment` in the `partitionSpec` can only be done during `PartialSegmentMergeTask` when the final merged segments are created.

3. Each worker receives the list of partitions ranges and for each firehose split, creates partial segments by using the partition ranges.

**NOTE:** In this initial version of implementing native parallel batch range partitioning, specifying the partition dimension is _required_ (unlike hadoop batch indexing).

### Rationale

Some other alternatives considered were:

#### Exact Histogram

Because an exact histogram would use much more memory than an approximate histogram (e.g., one from DataSketches), it would need to be periodically flushed to disk. To facilitate subsequent merging of histograms without consuming too much memory, the histogram would need to be serialized in sorted order so that a heap-merge can later be used. The main disadvantage of the exact histogram approach is reduced performance due to the amount of disk I/Os.

### Operational impact

None expected.

### Test plan

Unit tests and integrations tests, similar to the ones for native parallel batch indexing hashed partitioning.

### Future work

#### Automatic Partition Dimension Selection

This is a niche feature of hadoop batch indexing, and can be added later. Determining the best partition dimension can be complicated as many factors need to be considered (e.g., impact on query speed, impact on ingestion speed, etc.).

Contributor guide

Open the contributing guide

Research direction

Start at ParallelIndexSupervisorTask.runMultiPhaseParallel() and compare the existing PartialSegmentGenerateTask and PartialSegmentMergeTask flow for hashed partitioning. Read the proposed SingleDimensionPartitionsSpec behavior and the ItemsSketch/ItemsUnion approach, then review the existing native parallel hashed-partitioning unit and integration tests. Done means native parallel indexing accepts range partitions and the corresponding tests cover the new flow.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
data-engineering, distributed-systems
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
30/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.