Skip to content

[improvement](streaming) Support one-time S3 streaming ingestion - #68007

Open
JNSimba wants to merge 1 commit into
apache:masterfrom
JNSimba:feature/s3-once-streaming-ingestion
Open

JNSimba wants to merge 1 commit into
apache:masterfrom
JNSimba:feature/s3-once-streaming-ingestion

Conversation

@JNSimba

@JNSimba JNSimba commented Sep 15, 2026

Copy link
Copy Markdown
Member

What problem does this PR solve?

Issue Number: None

Related PR: None

Problem Summary: S3 streaming insert jobs currently keep polling for new files after consuming all files matched by the path. This change adds s3.ingestion_mode=ONCE, which imports matching files in lexical batches and marks the job as FINISHED after the final batch commits successfully. Recovered jobs probe from their committed offset and finish when no files remain. ONCE jobs reject user-specified offsets, including offset changes through ALTER JOB.

Release note

S3 streaming insert jobs support s3.ingestion_mode=ONCE for one-time, batched ingestion.

Check List (For Author)

  • Test

    • Unit Test
      • 19 tests passed: S3 offset batching and recovery, scheduler completion, task commit handling, and property validation.
    • Regression test
    • Manual test
    • No need to test or manual test.
  • Behavior changed:

    • Yes. S3 streaming jobs configured with ONCE finish after all discovered lexical batches are committed.
    • No.
  • Does this need documentation?

    • Yes. The new S3 ingestion mode needs to be added to the Streaming Job documentation.
    • No.

Check List (For Reviewer who merge this PR)

  • Confirm the release note
  • Confirm test cases
  • Confirm document
  • Add branch pick label

@hello-stephen

Copy link
Copy Markdown
Contributor

Thank you for your contribution to Apache Doris.
Don't know what should be done next? See How to process your PR.

Please clearly describe your PR:

  1. What problem was fixed (it's best to include specific error reporting information). How it was fixed.
  2. Which behaviors were modified. What was the previous behavior, what is it now, why was it modified, and what possible impacts might there be.
  3. What features were added. Why was this function added?
  4. Which code was refactored and why was this part of the code refactored?
  5. Which functions were optimized and what is the difference before and after the optimization?

@JNSimba
JNSimba force-pushed the feature/s3-once-streaming-ingestion branch from 7dc467a to 47a39f8 Compare September 15, 2026 08:09
@JNSimba

JNSimba commented Sep 15, 2026

Copy link
Copy Markdown
Member Author

run buildall

@JNSimba

JNSimba commented Sep 15, 2026

Copy link
Copy Markdown
Member Author

/review

@hello-stephen

Copy link
Copy Markdown
Contributor
TPC-H: Total hot run time: 17015 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpch-tools
Tpch sf100 test result on commit 47a39f8de697291c14092d7c2b2f333e8940196c, data reload: false

------ Round 1 ----------------------------------
============================================
q1	17639	3053	3045	3045
q2	2092	257	222	222
q3	10249	881	510	510
q4	4678	252	211	211
q5	7664	586	378	378
q6	134	114	94	94
q7	529	507	432	432
q8	9236	859	927	859
q9	3485	2411	2438	2411
q10	6504	856	720	720
q11	425	202	182	182
q12	627	266	206	206
q13	18097	1547	1157	1157
q14	166	155	137	137
q15	q16	426	396	372	372
q17	1415	899	856	856
q18	3037	2242	2269	2242
q19	1097	852	766	766
q20	377	282	206	206
q21	4852	1783	1868	1783
q22	338	270	226	226
Total cold run time: 93067 ms
Total hot run time: 17015 ms

----- Round 2, with runtime_filter_mode=off -----
============================================
q1	3387	3354	3358	3354
q2	500	391	374	374
q3	2249	2304	2215	2215
q4	1191	1180	893	893
q5	2198	2120	2131	2120
q6	174	122	86	86
q7	1026	960	804	804
q8	1577	1391	1393	1391
q9	3150	3111	3100	3100
q10	1859	1784	1629	1629
q11	353	270	253	253
q12	459	431	342	342
q13	1485	1536	1156	1156
q14	177	171	156	156
q15	q16	399	397	350	350
q17	3625	3263	3302	3263
q18	4865	4434	4757	4434
q19	873	818	962	818
q20	1024	979	834	834
q21	3835	3164	3357	3164
q22	419	365	341	341
Total cold run time: 34825 ms
Total hot run time: 31077 ms

@hello-stephen

Copy link
Copy Markdown
Contributor
TPC-DS: Total hot run time: 82202 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpcds-tools
TPC-DS sf100 test result on commit 47a39f8de697291c14092d7c2b2f333e8940196c, data reload: false

query5	4248	408	335	335
query6	377	137	125	125
query7	4976	428	231	231
query8	289	125	115	115
query9	8691	2898	2897	2897
query10	393	218	183	183
query11	5412	1044	932	932
query12	122	72	73	72
query13	1206	449	308	308
query14	6128	2215	2104	2104
query14_1	1988	1987	1961	1961
query15	179	121	111	111
query16	963	372	353	353
query17	804	453	371	371
query18	2335	343	249	249
query19	160	146	113	113
query20	72	70	69	69
query21	199	103	89	89
query22	5415	5353	5367	5353
query23	6680	6078	6062	6062
query23_1	6188	6029	6095	6029
query24	7303	1088	785	785
query24_1	775	775	815	775
query25	432	318	253	253
query26	1224	228	128	128
query27	2794	434	252	252
query28	4679	1520	1503	1503
query29	931	472	329	329
query30	254	159	128	128
query31	847	402	326	326
query32	134	71	72	71
query33	450	219	177	177
query34	978	837	473	473
query35	399	399	335	335
query36	566	557	528	528
query37	116	83	78	78
query38	1016	850	833	833
query39	493	492	490	490
query39_1	466	458	481	458
query40	204	87	76	76
query41	54	51	51	51
query42	73	70	71	70
query43	248	241	212	212
query44	1085	530	541	530
query45	108	103	98	98
query46	800	829	513	513
query47	773	769	735	735
query48	312	311	242	242
query49	533	258	209	209
query50	750	261	198	198
query51	8313	8145	8265	8145
query52	75	67	60	60
query53	192	206	150	150
query54	198	164	169	164
query55	71	62	54	54
query56	206	175	152	152
query57	682	751	664	664
query58	203	166	165	165
query59	1213	1229	1084	1084
query60	245	171	181	171
query61	127	138	114	114
query62	380	208	181	181
query63	184	136	134	134
query64	2662	746	539	539
query65	1613	1594	1587	1587
query66	1775	262	205	205
query67	10129	9696	9715	9696
query68	2738	1199	752	752
query69	341	247	202	202
query70	698	648	659	648
query71	245	183	166	166
query72	2304	1669	1482	1482
query73	644	605	344	344
query74	1572	1213	1123	1123
query75	1177	1095	950	950
query76	2286	705	519	519
query77	251	253	222	222
query78	3879	3746	3235	3235
query79	1518	808	578	578
query80	1194	312	266	266
query81	492	152	132	132
query82	900	126	95	95
query83	303	207	189	189
query84	303	106	88	88
query85	820	346	269	269
query86	391	174	170	170
query87	1031	971	900	900
query88	2770	2096	2131	2096
query89	289	196	174	174
query90	1949	130	130	130
query91	132	117	97	97
query92	81	66	71	66
query93	1408	1109	751	751
query94	652	237	193	193
query95	521	248	299	248
query96	794	579	262	262
query97	1075	1048	1018	1018
query98	147	139	133	133
query99	427	345	311	311
Total cold run time: 176596 ms
Total hot run time: 82202 ms

@hello-stephen

Copy link
Copy Markdown
Contributor
ClickBench: Total hot run time: 14.57 s
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/clickbench-tools
ClickBench test result on commit 47a39f8de697291c14092d7c2b2f333e8940196c, data reload: false

query1	0.01	0.00	0.01
query2	0.08	0.04	0.04
query3	0.25	0.11	0.11
query4	1.60	0.10	0.10
query5	0.17	0.15	0.16
query6	1.29	0.70	0.68
query7	0.03	0.01	0.01
query8	0.05	0.03	0.03
query9	0.30	0.21	0.21
query10	0.36	0.35	0.34
query11	0.15	0.12	0.11
query12	0.14	0.11	0.13
query13	0.32	0.31	0.31
query14	0.44	0.45	0.44
query15	0.35	0.35	0.34
query16	0.21	0.21	0.21
query17	0.72	0.67	0.66
query18	0.18	0.17	0.15
query19	1.16	1.19	1.16
query20	0.02	0.01	0.01
query21	15.45	0.14	0.12
query22	5.09	0.04	0.05
query23	16.18	0.25	0.10
query24	3.08	0.35	0.29
query25	0.11	0.04	0.03
query26	0.73	0.17	0.12
query27	0.04	0.02	0.03
query28	3.82	0.57	0.26
query29	12.46	3.20	2.55
query30	0.25	0.11	0.11
query31	2.75	0.36	0.17
query32	3.54	0.32	0.23
query33	1.48	1.34	1.55
query34	15.35	2.24	1.80
query35	1.76	1.75	1.70
query36	0.47	0.29	0.29
query37	0.08	0.04	0.04
query38	0.04	0.03	0.03
query39	0.04	0.02	0.02
query40	0.12	0.08	0.07
query41	0.08	0.03	0.02
query42	0.04	0.02	0.02
query43	0.04	0.03	0.03
Total cold run time: 90.83 s
Total hot run time: 14.57 s

@github-actions github-actions Bot 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.

Capped static review of exact head 47a39f8de697291c14092d7c2b2f333e8940196c against base 96d0ac68e84f92ed71acaad13560c076231a9266. Six verified findings remain (2 P1, 4 P2). The maximum third review round added the final settled-state finding, so this is truthfully capped/incomplete rather than a converged no-new-findings result.

Goal and scope: the stable ASCII happy path does batch matched S3 files lexically and reaches FINISHED after a successful final commit, and the change remains focused. The goal is not ready overall because fresh-empty lifecycle, terminal-state concurrency, source scoping, mixed-version leadership, cloud final-offset persistence, and recovered-task cleanup are incorrect.

Concurrency and lifecycle: transaction callback ordering, recurring metadata work, delayed task scheduling, manual pause/stop/resume, cancellation, and all three new completion sites were traced. FINISHED is not absorbing and eligibility-to-FINISHED transitions are not atomic; recovered completion also leaves a PENDING placeholder task and stale retry message.

Persistence and compatibility: current-version legacy images default to LEXICAL, provider reconstruction order is correct, the default offset-provider interface preserves other providers, and stable local/cloud committed offsets were traced. However, an old electable FE silently reinterprets a new ONCE job as lexical polling, and cloud recovery can checkpoint a FINISHED job with stale offsetProviderPersist. No FE-BE wire format, process config, or database storage format changes were introduced; the behavior-defining job property is persisted in the generic job map.

Configuration and parallel paths: ONCE is immutable through ALTER and rejects explicit offsets for S3, but validation occurs before source resolution, so CDC and FROM-source JDBC accept the S3-only setting as a silent no-op and can have valid offset configuration rejected.

Tests and other checkpoints: the changed unit/regression tests cover stable batching and uncontended completion, but omit end-to-end empty-source state, latch-based terminal races, non-S3 negative commands, mixed-version replay/failover, cloud recovery plus image round-trip, and recovered placeholder cleanup. No local build or test execution was permitted; the author reports 19 unit tests, while regression/manual testing are unchecked in the PR body. No material new performance issue was found in the bounded pagination path. Logging exists, but final offset/task observability is incorrect as noted. The new user-visible property still needs documentation, which the PR already acknowledges. No additional user review focus was provided.


Env.getCurrentEnv().getJobManager().getStreamingTaskManager().removeRunningTask(task);
if (offsetProvider.hasReachedEnd()) {
if (offsetProvider.hasReachedEnd(task.getRunningOffset())) {

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.

[P1] Make ONCE completion terminal across scheduler races

The commit thread holds the job lock through this block, but the recurring metadata task performs S3 I/O before acquiring it. If that listing fails while the final batch commits, this block sets FINISHED and removes the transaction callback; the waiting fetchMeta path then changes the job to PAUSED because AbstractJob permits FINISHED -> PAUSED. Auto/manual resume can reopen a job whose callback is no longer registered. Both new scheduler completion paths have the inverse TOCTOU: their earlier eligibility checks are released before an unconditional FINISHED write, so they can overwrite a concurrent PAUSED/STOPPED state. Please make these transitions conditional under the same job lock, or enforce that terminal states are absorbing, and add latch-based coverage.

public static SourceOffsetProvider createSourceOffsetProvider(
String sourceType, StreamingJobProperties jobProperties) {
try {
if ("s3".equalsIgnoreCase(sourceType) && jobProperties.isS3OnceMode()) {

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.

[P1] Fence ONCE jobs from older FE leaders

The mode is persisted only as a key in the generic properties map. An FE at the base version replays that map without validating unknown keys, then its one-argument factory always reconstructs the no-arg lexical S3 provider. During a rolling FE upgrade, an older follower that becomes leader can therefore turn a PENDING ONCE job into indefinite polling and ingest later files instead of reaching FINISHED. Please gate creation until every electable FE supports this mode, or persist a backward-compatible discriminator that makes old leaders fail closed, with a mixed-version replay/failover test.

GlobListing globListing = fileSystem.globListWithLimit(Location.of(filePath), startFile, 1, 1);
if (!globListing.getFiles().isEmpty() && StringUtils.isNotEmpty(globListing.getMaxFile())) {
boolean hasFiles = !globListing.getFiles().isEmpty();
if (onceMode && startFile != null && !hasFiles) {

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.

[P2] Keep an initially empty ONCE source out of the failure budget

With no committed offset, this branch deliberately leaves reachedEnd false, while hasMoreDataToConsume returns true for currentOffset == null. The PENDING path does not probe metadata, so a normal empty glob is dispatched; getNextOffset throws, the task retries four times, and repeated auto-resumes can end in CANNOT_RESUME_ERR. ONCE should either finish an empty snapshot or wait through the no-data delay path, but it should not treat zero matches as a task failure. Please cover the full PENDING -> scheduler -> task state transition, not only the raw provider exception.

&& v <= (long) (1024 * 1024 * 1024) * 10,
StreamingJobProperties.S3_MAX_BATCH_BYTES_PROPERTY + " should between 100MB and 10GB");

String ingestionMode = getS3IngestionMode();

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.

[P2] Reject the S3-only mode on non-S3 sources

This validation runs before the source type is resolved, so cdc_stream and FROM-source JDBC jobs accept s3.ingestion_mode=ONCE even though their provider paths ignore it and continue indefinitely. The same irrelevant property can still reject those providers' supported offset setting here. Please validate the explicit property after source/provider resolution in CREATE and ALTER, rejecting it outside S3, and add negative coverage for both non-S3 paths.

}
// reject task if no more data to consume
if (!job.hasMoreDataToConsume()) {
if (job.hasReachedEnd()) {

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.

[P2] Persist the cloud-replayed offset before finishing

After a crash following a cloud commit, replayOnCloudMode restores the authoritative offset only into offsetProvider; unlike local replay, it does not refresh offsetProviderPersist. This new exhaustion branch can then journal FINISHED with the stale/null persisted string. Because Gson excludes the provider object and a restored FINISHED job does not replay Meta Service again, a later checkpoint restart permanently shows the old final offset. Please synchronize offsetProviderPersist whenever cloud replay applies an attachment before allowing either natural-completion path to persist FINISHED, and add a cloud recovery plus image-round-trip test.

if (job.hasReachedEnd()) {
job.updateJobStatus(JobStatus.FINISHED);
job.logUpdateOperation();
return;

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.

[P2] Retire the placeholder task when recovery finishes

A recovered ONCE job records a PENDING placeholder task; its first scheduler pass delays that task and sets the "will retry" runtime message. After the metadata probe marks the source exhausted, the delayed pass reaches this return without executing, canceling, finalizing, or removing the task and without clearing the message. The job is therefore shown as FINISHED while its task remains PENDING with no finish time and the job still says it will retry. Please finalize/cancel the placeholder and clear its reference/message as part of the terminal transition, with a recovered-empty-tail state test.

@hello-stephen

Copy link
Copy Markdown
Contributor

FE UT Coverage Report

Increment line coverage 72.92% (35/48) 🎉
Increment coverage report
Complete coverage report

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.

2 participants