Skip to content

Admit verified at-least-once Paimon consumer retention - #287

Merged
jordepic merged 1 commit into
mainfrom
work/paimon-consumer-retention
Sep 29, 2026
Merged

jordepic merged 1 commit into
mainfrom
work/paimon-consumer-retention

Conversation

@jordepic

Copy link
Copy Markdown
Collaborator

Explicit Paimon 2.0 consumer.mode=at-least-once sources with a configured expiration can now use native Parquet/ORC decoding. Previously every consumer ID forced the stock source, although the native reader already used Paimon's continuous enumerator, progress events and checkpoint serializers.

Part of #27; this does not close the broader source-coverage issue. Default/exactly-once consumers retain Paimon's dedicated host topology. Other explicit consumer options remain gated. There are no JNI, runtime-reader or checkpoint-format changes.

Validation:

  • 25 focused admission, reader recovery and consumer-retention tests passed in each of four combinations: released Flink 2.2.1/1.18.1 with Parquet/ORC (100 executions).
  • Independent host/native jobs verify durable cursor resumption. Coordinator tests verify only completed checkpoints advance retention, serialized enumerator restoration, protection from snapshot expiration and release after progress. Reader tests retain within-batch offsets across snapshot/tail and merge modes.
  • Release+mimalloc benchmarks use two warmups/five trials, rotating order, parallelism 1, 2 GiB heap, 100 ms checkpoints, nullable strings/BIGINTs, 200K and 2M rows, native file decoding checks and actual substitution assertions. They include planning, job startup, reads, JNI and row conversions.
  • At 2M rows, median whole-job time falls 27–33% versus stock and 26–32% versus the previous source route. “Previous” is a source-disable ablation on base 665dc10d, not a historical binary rebuild. Full ranges and every trial are in the connector docs and CSV.
  • Append uses SQL blackhole. Primary-key uses a full-changelog rowwise DataStream discard sink: SQL blackhole inserts DropUpdateBefore and still triggers the existing whole-island fallback. These bounded insert fixtures do not establish sustained updating-stream throughput.
  • Retained the unfavorable 200K pilot (native 0.162 s versus previous 0.115 s); its per-native-job EXPLAIN created asymmetric planner warming. Final measurements validate plans independently before all modes. Small jobs remain startup-sensitive.
  • Strict MkDocs, native package/JNI boundary check and diff whitespace checks passed. Final benchmark source also compiles on Flink 1.18.1.

Coverage, reproduction command, sink limitations, variability and raw results are documented in docs/connectors/paimon.md and docs/benchmarks/paimon-consumer-2026-09-28.csv.

Paimon's continuous enumerator already owns checkpointed consumer progress,
and our source retains its split events and checkpoint serializers. Permit
explicit at-least-once consumers with expiration configured after proving
completed-checkpoint cursor advancement, restored reader offsets, fresh-job
resumption and snapshot protection/expiration against the released source.
Keep default/exactly-once consumers on their dedicated host topology and
retain fallback for other unverified consumer options.

Release+mimalloc whole-job comparisons on Flink 2.2.1 reduce median time
27-33% versus stock and 26-32% versus the prior host-source route at 2M rows
across Parquet/ORC append and primary-key fixtures. Retain row conversions,
all trials, startup variability and the unfavorable asymmetric-warmup pilot.
Primary-key SQL blackhole still falls back at DropUpdateBefore; its measured
control uses a full-changelog rowwise DataStream sink, documented explicitly.

Validate 25 focused tests on each Flink 2.2.1/1.18.1 and Parquet/ORC
combination, final 1.18 benchmark compilation, native module boundaries,
and strict MkDocs. No native runtime or checkpoint format changes.

Part of #27.
@jordepic
jordepic merged commit 2f66ab8 into main Sep 29, 2026
48 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant