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

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


The following commit(s) were added to refs/heads/master by this push:
     new 6fe144b3768 [BEAM-40264] Fix Spark MapState and SetState isEmpty after 
removals (#40398)
6fe144b3768 is described below

commit 6fe144b376897b58cd8facfa563e976a128d74b3
Author: Dhruv Dagar <[email protected]>
AuthorDate: Tue Oct 6 23:27:05 2026 +0530

    [BEAM-40264] Fix Spark MapState and SetState isEmpty after removals (#40398)
    
    * [Go SDK] Add portable logical type abstraction
    
    * Fix Spark MapState and SetState emptiness after removals
    
    * Add regression tests for state emptiness after removal
    
    * Remove unrelated upstream file from BEAM-40264 branch
---
 .../org/apache/beam/runners/core/StateInternalsTest.java     | 10 ++++++++++
 .../beam/runners/spark/stateful/SparkStateInternals.java     | 12 ++++++++++--
 2 files changed, 20 insertions(+), 2 deletions(-)

diff --git 
a/runners/core-java/src/test/java/org/apache/beam/runners/core/StateInternalsTest.java
 
b/runners/core-java/src/test/java/org/apache/beam/runners/core/StateInternalsTest.java
index e15249969f2..6988a492736 100644
--- 
a/runners/core-java/src/test/java/org/apache/beam/runners/core/StateInternalsTest.java
+++ 
b/runners/core-java/src/test/java/org/apache/beam/runners/core/StateInternalsTest.java
@@ -225,6 +225,11 @@ public abstract class StateInternalsTest {
     assertThat(later.read(), hasItems("C", "D"));
     assertFalse(later.contains("A").read());
 
+    value.remove("B");
+    value.remove("C");
+    value.remove("D");
+    assertTrue(value.isEmpty().read());
+
     // clear
     value.clear();
     assertThat(value.read(), Matchers.emptyIterable());
@@ -389,6 +394,11 @@ public abstract class StateInternalsTest {
     // isEmpty
     assertFalse(value.isEmpty().read());
 
+    value.remove("B");
+    value.remove("D");
+    value.remove("E");
+    assertTrue(value.isEmpty().read());
+
     // clear
     value.clear();
     assertThat(value.entries().read(), Matchers.emptyIterable());
diff --git 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkStateInternals.java
 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkStateInternals.java
index 51ceb4c8730..4f744ab3ab1 100644
--- 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkStateInternals.java
+++ 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkStateInternals.java
@@ -432,7 +432,11 @@ public class SparkStateInternals<K> implements 
StateInternals {
     public void remove(MapKeyT key) {
       Map<MapKeyT, MapValueT> sparkMapState = readAsMap();
       sparkMapState.remove(key);
-      writeValue(sparkMapState);
+      if (sparkMapState.isEmpty()) {
+        clear();
+      } else {
+        writeValue(sparkMapState);
+      }
     }
 
     @Override
@@ -537,7 +541,11 @@ public class SparkStateInternals<K> implements 
StateInternals {
     public void remove(InputT input) {
       Set<InputT> sparkSetState = readAsSet();
       sparkSetState.remove(input);
-      writeValue(sparkSetState);
+      if (sparkSetState.isEmpty()) {
+        clear();
+      } else {
+        writeValue(sparkSetState);
+      }
     }
 
     @Override

Reply via email to