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