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();
}