GAURAV MIGLANI
05/11/2026, 9:23 AMTaranpreet Kaur
05/13/2026, 5:45 AMAlbert Battalov
05/19/2026, 8:44 PMpublic class JitteredEventTimeTrigger extends Trigger<Object, TimeWindow> {
private final long maxJitterMillis;
@Override
public TriggerResult onElement(Object element, long timestamp, TimeWindow window, TriggerContext ctx) {
ctx.registerEventTimeTimer(window.maxTimestamp());
return TriggerResult.CONTINUE;
}
@Override
public TriggerResult onEventTime(long time, TimeWindow window, TriggerContext ctx) {
long jitter = ThreadLocalRandom.current().nextLong(maxJitterMillis);
ctx.registerProcessingTimeTimer(ctx.getCurrentProcessingTime() + jitter);
return TriggerResult.CONTINUE;
}
@Override
public TriggerResult onProcessingTime(long time, TimeWindow window, TriggerContext ctx) {
return TriggerResult.FIRE;
}
@Override
public void clear(TimeWindow window, TriggerContext ctx) {
ctx.deleteEventTimeTimer(window.maxTimestamp());
}
}
Not seeing any output. I think WindowOperator https://github.com/apache/flink/blob/master/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java#L485-L487 might be cleaning up window state at maxTimestamp + allowedLateness(0) in the same onEventTime call - before my processing-time timer gets a chance to fire. Can anyone confirm?
Flink 2.2 running on AWS Managed FlinkEvan NotEvan
05/21/2026, 1:41 PMFlink SQL> SET 'execution.checkpointing.interval' = 0;
[ERROR] Could not execute SQL statement. Reason:
org.apache.flink.sql.parser.impl.ParseException: Encountered "0" at line 1, column 42.
Was expecting one of:
<BINARY_STRING_LITERAL> ...
<QUOTED_STRING> ...
<PREFIXED_STRING_LITERAL> ...
<UNICODE_STRING_LITERAL> ...
<C_STYLE_ESCAPED_STRING_LITERAL> ...
<BIG_QUERY_DOUBLE_QUOTED_STRING> ...
<BIG_QUERY_QUOTED_STRING> ...
Flink SQL> SET 'execution.checkpointing.interval' = 'false';
[ERROR] Could not execute SQL statement. Reason:
java.lang.NumberFormatException: text does not start with a number, and is not a valid ISO-8601 duration format: false
Flink SQL> SET 'execution.checkpointing.interval' = '0s';
[ERROR] Could not execute SQL statement. Reason:
java.lang.IllegalArgumentException: Checkpoint interval must be larger than or equal to 10 ms
Flink SQL> SET 'execution.checkpointing.interval' = '0 s';
[ERROR] Could not execute SQL statement. Reason:
java.lang.IllegalArgumentException: Checkpoint interval must be larger than or equal to 10 ms
Flink SQL> SET 'execution.checkpointing.interval' = NONE;
[ERROR] Could not execute SQL statement. Reason:
org.apache.flink.sql.parser.impl.ParseException: Encountered "NONE" at line 1, column 42.
Was expecting one of:
<BINARY_STRING_LITERAL> ...
<QUOTED_STRING> ...
<PREFIXED_STRING_LITERAL> ...
<UNICODE_STRING_LITERAL> ...
<C_STYLE_ESCAPED_STRING_LITERAL> ...
<BIG_QUERY_DOUBLE_QUOTED_STRING> ...
<BIG_QUERY_QUOTED_STRING> ...
Flink SQL> SET 'execution.checkpointing.interval' = NULL;
[ERROR] Could not execute SQL statement. Reason:
org.apache.flink.sql.parser.impl.ParseException: Encountered "NULL" at line 1, column 42.
Was expecting one of:
<BINARY_STRING_LITERAL> ...
<QUOTED_STRING> ...
<PREFIXED_STRING_LITERAL> ...
<UNICODE_STRING_LITERAL> ...
<C_STYLE_ESCAPED_STRING_LITERAL> ...
<BIG_QUERY_DOUBLE_QUOTED_STRING> ...
<BIG_QUERY_QUOTED_STRING> ...Albert Battalov
05/26/2026, 3:50 AMFabrizzio Chavez
05/27/2026, 6:05 AMAnish Kulkarni
05/27/2026, 8:29 AMCannot map checkpoint/savepoint state for operator f17fc775f6e82bc8d571beb78ac6c693 to the new program, because the operator is not available in the new program. If you want to allow to skip this, you can set the --allowNonRestoredState option on the CLI.
I want to avoid the usage of allowNonRestoredState flag since it could hide real issues
What is the way to achieve starting from snapshots safely after a change in operator chain/sources?Dominik
05/29/2026, 8:37 AMFrancis Altomare
06/02/2026, 2:43 PMupgradeMode last-state and using the adaptive scheduler we observe:
• Downscale rescales (e.g. parallelism 24 -> 12) take roughly 2 hours from spec change to RUNNING (no events are processed in this time
• Upscale rescales in the same flow are much faster, taking around 12 minutes
The relevant ForSt / checkpointing settings are:
• use-ingest-db-restore-mode: true
• incremental-restore-async-compact-after-rescale: true
• state-recovery.claim-mode: CLAIM
• NATIVE savepoint format, incremental + compressed checkpoints
• pipeline.max-parallelism: 120
• ForSt cache size-based-limit 350GB on a 500Gi gp3 volume per TM.
• I’m running one slot per TM and have tried parallelisms between 12 and 24
My questions are:
• Is ~2h for a downscale at this state size in line with what others see, or does it suggest a misconfiguration? Most of the time appears to be in restore.
• Is the asymmetry where upscales are much faster also expected?
• Given the asymmetry, are there settings or patterns that make autoscaling viable at ~1TB state beyond adjusting job.autoscaler.scale-down.max-factor and scale-down.interval? Or is the practical recommendation at this scale to keep parallelism static and rescale manually during low-traffic windows?
Thank you!Bianca Falcone
06/03/2026, 8:48 AMDanny Wilkins
06/08/2026, 3:52 PMpekko tests try to connect to my machine's local-but-external ip (ie 10.0.0.whatever.) Is there a property or something else I can just -D in order to force pekko to use localhost instead of whatever discovery mechanism it's doing?HunkSurvivor89
06/08/2026, 6:07 PMHan You
06/23/2026, 6:08 PMOr Keren
07/01/2026, 12:30 PMAsyncException{org.forstdb.RocksDBException: Exception when Flush file, path: <s3://yotpo-cdp/flinkcdp/checkpoints/shared/op_AsyncKeyedProcessOperator_f199ad3ac3e86e5fc9372ebb239a06cb__33_64__attempt_0/db/260675.sst>}
at org.apache.flink.runtime.asyncprocessing.operators.AbstractAsyncKeyOrderedStreamOperator.handleAsyncException(AbstractAsyncKeyOrderedStreamOperator.java:120)
at org.apache.flink.core.asyncprocessing.AsyncFutureImpl.completeExceptionally(AsyncFutureImpl.java:553)
at org.apache.flink.state.forst.ForStDBPutRequest.completeStateFutureExceptionally(ForStDBPutRequest.java:87)
at org.apache.flink.state.forst.ForStWriteBatchOperation.lambda$process$0(ForStWriteBatchOperation.java:80)
at java.base/java.util.concurrent.CompletableFuture$AsyncRun.run(Unknown Source)
at org.apache.flink.util.concurrent.DirectExecutorService.execute(DirectExecutorService.java:247)
at java.base/java.util.concurrent.CompletableFuture.asyncRunStage(Unknown Source)
at java.base/java.util.concurrent.CompletableFuture.runAsync(Unknown Source)
at org.apache.flink.state.forst.ForStWriteBatchOperation.process(ForStWriteBatchOperation.java:65)
at org.apache.flink.state.forst.ForStStateExecutor.lambda$executeBatchRequests$2(ForStStateExecutor.java:201)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source)
at org.apache.flink.state.forst.ForStStateExecutor$ForStExecutorThreadFactory.lambda$newThread$0(ForStStateExecutor.java:338)
at java.base/java.lang.Thread.run(Unknown Source)
Caused by: org.forstdb.RocksDBException: Exception when Flush file, path: <s3://yotpo-cdp/flinkcdp/checkpoints/shared/op_AsyncKeyedProcessOperator_f199ad3ac3e86e5fc9372ebb239a06cb__33_64__attempt_0/db/260675.sst>
at org.forstdb.RocksDB.write0(Native Method)
at org.forstdb.RocksDB.write(RocksDB.java:1843)
at org.apache.flink.state.forst.ForStDBWriteBatchWrapper.flush(ForStDBWriteBatchWrapper.java:136)
at org.apache.flink.state.forst.ForStDBWriteBatchWrapper.flushIfNeeded(ForStDBWriteBatchWrapper.java:155)
at org.apache.flink.state.forst.ForStDBWriteBatchWrapper.put(ForStDBWriteBatchWrapper.java:117)
at org.apache.flink.state.forst.ForStDBPutRequest.process(ForStDBPutRequest.java:68)
at org.apache.flink.state.forst.ForStWriteBatchOperation.lambda$process$0(ForStWriteBatchOperation.java:70)
... 10 more
Suppressed: org.forstdb.RocksDBException: Exception when Flush file, path: <s3://yotpo-cdp/flinkcdp/checkpoints/shared/op_AsyncKeyedProcessOperator_f199ad3ac3e86e5fc9372ebb239a06cb__33_64__attempt_0/db/260675.sst>
at org.forstdb.RocksDB.write0(Native Method)
at org.forstdb.RocksDB.write(RocksDB.java:1843)
at org.apache.flink.state.forst.ForStDBWriteBatchWrapper.flush(ForStDBWriteBatchWrapper.java:136)
at org.apache.flink.state.forst.ForStDBWriteBatchWrapper.close(ForStDBWriteBatchWrapper.java:144)
at org.apache.flink.state.forst.ForStWriteBatchOperation.lambda$process$0(ForStWriteBatchOperation.java:67)
... 10 moreFabrizzio Chavez
07/03/2026, 1:13 AMSơn Bùi
07/08/2026, 4:09 AMKaran Kumar
07/09/2026, 5:25 AMDavid Causse
07/10/2026, 2:07 PMcompatibleAfterMigration. Initially my understanding was that the old serializer is used to read the stored bytes, the instance object is then written with the new serializer and these new bytes are read by the new serializer, in short I would have assumed that it ensured that the instances received in the operators are always constructed using the new serializer. But I'm hitting a bug where I suspect that the instance read by the old serializer is simply passed to the operator without any migration. (Note this is the operator state of the Async operator). Could someone clarify my understanding about state migration with compatibleAfterMigration? Thanks! 🙏pavithranagayallapu
07/17/2026, 11:12 AMMaitreya Manohar
07/20/2026, 7:05 AMpavithranagayallapu
07/20/2026, 11:05 AMzhongxia.zhou
07/20/2026, 9:20 PMFrancis Altomare
07/22/2026, 8:22 AMpavithranagayallapu
07/23/2026, 3:31 AMpavithranagayallapu
07/23/2026, 3:43 AMSiddhesh Kalgaonkar
08/04/2026, 3:43 PMUdith V B
08/06/2026, 1:21 PMNishtha Bhattacharjee
08/09/2026, 9:19 PMnoa.cavassi
08/12/2026, 12:20 PM[ERROR] Could not execute SQL statement. Reason:
java.lang.ClassNotFoundException: org.apache.iceberg.data.IcebergGenericReader
... 15 more
And the full error in the jobmanager out file:
java.lang.RuntimeException: One or more fetchers have encountered exception
at org.apache.flink.connector.base.source.reader.fetcher.SplitFetcherManager.checkErrors(SplitFetcherManager.java:333) ~[flink-connector-files-1.20.0.jar:1.20.0]
at org.apache.flink.connector.base.source.reader.SourceReaderBase.getNextFetch(SourceReaderBase.java:228) ~[flink-connector-files-1.20.0.jar:1.20.0]
at org.apache.flink.connector.base.source.reader.SourceReaderBase.pollNext(SourceReaderBase.java:190) ~[flink-connector-files-1.20.0.jar:1.20.0]
at org.apache.flink.streaming.api.operators.SourceOperator.emitNext(SourceOperator.java:422) ~[flink-dist-1.20.0.jar:1.20.0]
at org.apache.flink.streaming.runtime.io.StreamTaskSourceInput.emitNext(StreamTaskSourceInput.java:68) ~[flink-dist-1.20.0.jar:1.20.0]
at org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:65) ~[flink-dist-1.20.0.jar:1.20.0]
at org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:638) ~[flink-dist-1.20.0.jar:1.20.0]
at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:231) ~[flink-dist-1.20.0.jar:1.20.0]
at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:973) ~[flink-dist-1.20.0.jar:1.20.0]
at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:917) ~[flink-dist-1.20.0.jar:1.20.0]
at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:970) ~[flink-dist-1.20.0.jar:1.20.0]
at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:949) ~[flink-dist-1.20.0.jar:1.20.0]
at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:763) ~[flink-dist-1.20.0.jar:1.20.0]
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:575) ~[flink-dist-1.20.0.jar:1.20.0]
at java.lang.Thread.run(Thread.java:1583) ~[?:?]
Caused by: java.lang.NoClassDefFoundError: org/apache/iceberg/data/IcebergGenericReader
at org.apache.fluss.lake.iceberg.source.IcebergRecordReader.<init>(IcebergRecordReader.java:65) ~[fluss-lake-iceberg-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.fluss.lake.iceberg.source.IcebergLakeSource.createRecordReader(IcebergLakeSource.java:100) ~[fluss-lake-iceberg-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.fluss.flink.lake.reader.LakeSnapshotScanner.pollBatch(LakeSnapshotScanner.java:54) ~[fluss-flink-1.20-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.fluss.flink.source.reader.BoundedSplitReader.pollBatch(BoundedSplitReader.java:128) ~[fluss-flink-1.20-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.fluss.flink.source.reader.BoundedSplitReader.poll(BoundedSplitReader.java:121) ~[fluss-flink-1.20-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.fluss.flink.source.reader.BoundedSplitReader.readBatch(BoundedSplitReader.java:76) ~[fluss-flink-1.20-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.fluss.flink.source.reader.FlinkSourceSplitReader.fetch(FlinkSourceSplitReader.java:154) ~[fluss-flink-1.20-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.flink.connector.base.source.reader.fetcher.FetchTask.run(FetchTask.java:58) ~[flink-connector-files-1.20.0.jar:1.20.0]
at org.apache.flink.connector.base.source.reader.fetcher.SplitFetcher.runOnce(SplitFetcher.java:165) ~[flink-connector-files-1.20.0.jar:1.20.0]
at org.apache.flink.connector.base.source.reader.fetcher.SplitFetcher.run(SplitFetcher.java:117) ~[flink-connector-files-1.20.0.jar:1.20.0]
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:572) ~[?:?]
at java.util.concurrent.FutureTask.run(FutureTask.java:317) ~[?:?]
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1144) ~[?:?]
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:642) ~[?:?]
... 1 more
Caused by: java.lang.ClassNotFoundException: org.apache.iceberg.data.IcebergGenericReader
at jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:641) ~[?:?]
at jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(ClassLoaders.java:188) ~[?:?]
at java.lang.ClassLoader.loadClass(ClassLoader.java:526) ~[?:?]
at org.apache.fluss.lake.iceberg.source.IcebergRecordReader.<init>(IcebergRecordReader.java:65) ~[fluss-lake-iceberg-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.fluss.lake.iceberg.source.IcebergLakeSource.createRecordReader(IcebergLakeSource.java:100) ~[fluss-lake-iceberg-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.fluss.flink.lake.reader.LakeSnapshotScanner.pollBatch(LakeSnapshotScanner.java:54) ~[fluss-flink-1.20-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.fluss.flink.source.reader.BoundedSplitReader.pollBatch(BoundedSplitReader.java:128) ~[fluss-flink-1.20-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.fluss.flink.source.reader.BoundedSplitReader.poll(BoundedSplitReader.java:121) ~[fluss-flink-1.20-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.fluss.flink.source.reader.BoundedSplitReader.readBatch(BoundedSplitReader.java:76) ~[fluss-flink-1.20-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.fluss.flink.source.reader.FlinkSourceSplitReader.fetch(FlinkSourceSplitReader.java:154) ~[fluss-flink-1.20-0.9.1-incubating.jar:0.9.1-incubating]
at org.apache.flink.connector.base.source.reader.fetcher.FetchTask.run(FetchTask.java:58) ~[flink-connector-files-1.20.0.jar:1.20.0]
at org.apache.flink.connector.base.source.reader.fetcher.SplitFetcher.runOnce(SplitFetcher.java:165) ~[flink-connector-files-1.20.0.jar:1.20.0]
at org.apache.flink.connector.base.source.reader.fetcher.SplitFetcher.run(SplitFetcher.java:117) ~[flink-connector-files-1.20.0.jar:1.20.0]
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:572) ~[?:?]
at java.util.concurrent.FutureTask.run(FutureTask.java:317) ~[?:?]
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1144) ~[?:?]
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:642) ~[?:?]
... 1 more
Any ideas on what the problem could be? I'm happy to provide any other information if necessary.Gwenael LE BARZIC
08/31/2026, 3:40 PM