This is an automated email from the ASF dual-hosted git repository.

gustavodemorais pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git


The following commit(s) were added to refs/heads/master by this push:
     new 79b4843ac25 [FLINK-40339][table] Do not read buffer entries after 
removing them
79b4843ac25 is described below

commit 79b4843ac2579f97387d18be58254a849f437866
Author: Gustavo de Morais <[email protected]>
AuthorDate: Thu Aug 6 17:56:47 2026 +0200

    [FLINK-40339][table] Do not read buffer entries after removing them
    
    This closes #28933.
---
 .../runtime/operators/sink/WatermarkCompactingSinkMaterializer.java   | 4 +++-
 1 file changed, 3 insertions(+), 1 deletion(-)

diff --git 
a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/WatermarkCompactingSinkMaterializer.java
 
b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/WatermarkCompactingSinkMaterializer.java
index 40624d64389..6e6e1c0e51c 100644
--- 
a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/WatermarkCompactingSinkMaterializer.java
+++ 
b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/WatermarkCompactingSinkMaterializer.java
@@ -219,7 +219,9 @@ public class WatermarkCompactingSinkMaterializer extends 
TableStreamOperator<Row
             List<Map.Entry<Long, List<RowData>>> entries = new ArrayList<>();
             Iterator<Map.Entry<Long, List<RowData>>> iterator = 
buffer.entries().iterator();
             while (iterator.hasNext()) {
-                entries.add(iterator.next());
+                // The value of an entry is undefined once the iterator 
removed it.
+                final Map.Entry<Long, List<RowData>> entry = iterator.next();
+                entries.add(Map.entry(entry.getKey(), entry.getValue()));
                 iterator.remove();
             }
 

Reply via email to