apache / apache/pinot

The Pinot Flink connector does not gracefully handle failure

Open
#8,889 3 comments 0 reactions 0 assignees View on GitHub
Dominant language
Java
Stars
6.1k
Forks
1.5k
Avg merge
1d 21h
Merged PRs (30d)
189

Description

The current pinot flink connector does not gracefully handle errors. Due to the way the connector works, if it errors in the middle of adding segments to a table, the table ends up with an inconsistent view. Additionally, the connector does not currently support refresh tables. Refresh tables require atomic segment replacement, but the connector currently naively uploads segments as they are built.

From testing the connector in production, I've also identified a few performance issues. These have a few different causes; The AVRO serialization is not configurable, nor is the file writing configurable (for example for different block sizes).

I have written a flink connector based on this one, but with some heavy amendments. First of all, it implements WithPostCommitTopology from flink, implementing a global committer. It does work in a few different stages:

1. Operator is responsible for sending serialized AVRO records directly to the sink
2. The sink writer is responsible for building and flushing segments to disk
3. The sink committer (before global commit) is responsible for uploading the segments to a location that is reachable by all nodes in the flink cluster (In my case, to S3 deep store)
4. The global sink committer executes the segment replacement protocol defined in the Pinot SDK.

This sink is currently only compatible with REFRESH type tables that want to replace all segments on every single job execution. It takes care of atomically replacing the segments for the table, and performs well due to the way it does the hard work upfront. I am open to sharing this code so that it can be merged into the pinot repository, but it does have some limitations.

- No checkpointing
- Only BATCH execution mode is supported at the moment
- Only REFRESH tables are supported at the moment (Full segment replacement)
- The connector currently bypasses certain Pinot conventions (such as using certain attributes defined in the batch config, and so on). This would need to be approached with scrutiny to ensure the code is in-line with the rest of the repository.

Contributor guide

Open the contributing guide

Research direction

Start by reviewing the existing Pinot Flink connector and the proposed four-stage sink design described in the issue. Compare the failure handling, REFRESH-table support, serialization and file-writing configuration, checkpointing, and execution-mode limitations. Done requires an agreed scope and implementation plan before code can be started.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
data-engineering, databases
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
20/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.