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

stankiewicz 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 76a53714d2b Nullness improvements for simple core transforms (#40022)
76a53714d2b is described below

commit 76a53714d2b4fc38071d2a325132af2cee5e501d
Author: Kenneth Knowles <[email protected]>
AuthorDate: Tue Sep 29 09:44:33 2026 -0400

    Nullness improvements for simple core transforms (#40022)
    
    * Mark TimestampedValue as covariant in its value type
    
    TimestampedValue is immutable, so a TimestampedValue<V> is safe to use
    wherever a TimestampedValue of a supertype of V is expected. Saying so
    lets the nullness checker accept passing a TimestampedValue<T> where a
    TimestampedValue<@Nullable T> is wanted, which combiners need when their
    accumulator starts out holding a null value.
    
    This only loosens checking, so it cannot invalidate existing callers.
    
    * Nullness improvements for Latest
    
    Replace the blanket class-level nullness suppression with narrow,
    documented suppressions on the methods that genuinely produce a null.
    
    LatestFn's accumulator holds a null value until the first input arrives,
    so it is now typed TimestampedValue<@Nullable T>. That is invisible in
    the public API, since Latest.combineFn() erases AccumT to a wildcard.
    It relies on TimestampedValue being covariant in its value type, so that
    an input can still be passed where an accumulator is expected.
    
    No runtime changes beyond hoisting the timestamp comparison into a
    private helper so that mergeAccumulators no longer routes through
    addInput. The null-element check it relied on is retained.
    
    * Nullness improvements for Distinct
    
    Remove the blanket class-level nullness suppression. The representative
    type of WithRepresentativeValues is genuinely optional -- expand()
    already branches on it -- so annotate it rather than suppressing.
    
    No runtime changes.
    
    * Nullness improvements for Reify
    
    Remove the blanket class-level nullness suppression. Reifying a view in
    the global window creates a single placeholder element, which is a null
    Void; spell that type out rather than casting.
    
    No runtime changes.
    
    * Nullness improvements for Sample
    
    Remove the blanket class-level nullness suppression, keeping the
    unrelated rawtypes one. AnyValueCombineFn.extractOutput returns null
    when the input is empty, and OutputT is T in the public signature of
    anyValueCombineFn(), so that keeps a narrow suppression.
    
    No runtime changes.
---
 .../org/apache/beam/sdk/transforms/Distinct.java   |  7 +--
 .../org/apache/beam/sdk/transforms/Latest.java     | 52 ++++++++++++++--------
 .../java/org/apache/beam/sdk/transforms/Reify.java |  6 +--
 .../org/apache/beam/sdk/transforms/Sample.java     |  6 +--
 .../apache/beam/sdk/values/TimestampedValue.java   |  3 ++
 .../apache/beam/sdk/transforms/LatestFnTest.java   |  6 +--
 6 files changed, 46 insertions(+), 34 deletions(-)

diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Distinct.java 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Distinct.java
index 641412772bf..852e1fb31c9 100644
--- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Distinct.java
+++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Distinct.java
@@ -61,9 +61,6 @@ import org.joda.time.Duration;
  *
  * @param <T> the type of the elements of the input and output {@code 
PCollection}s
  */
-@SuppressWarnings({
-  "nullness" // TODO(https://github.com/apache/beam/issues/20497)
-})
 public class Distinct<T> extends PTransform<PCollection<T>, PCollection<T>> {
 
   /**
@@ -148,10 +145,10 @@ public class Distinct<T> extends 
PTransform<PCollection<T>, PCollection<T>> {
   public static class WithRepresentativeValues<T, IdT>
       extends PTransform<PCollection<T>, PCollection<T>> {
     private final SerializableFunction<T, IdT> fn;
-    private final TypeDescriptor<IdT> representativeType;
+    private final @Nullable TypeDescriptor<IdT> representativeType;
 
     private WithRepresentativeValues(
-        SerializableFunction<T, IdT> fn, TypeDescriptor<IdT> 
representativeType) {
+        SerializableFunction<T, IdT> fn, @Nullable TypeDescriptor<IdT> 
representativeType) {
       this.fn = fn;
       this.representativeType = representativeType;
     }
diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Latest.java 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Latest.java
index cd52f05cad5..dc75ea9d3b0 100644
--- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Latest.java
+++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Latest.java
@@ -17,8 +17,8 @@
  */
 package org.apache.beam.sdk.transforms;
 
+import static org.apache.beam.sdk.util.Preconditions.checkArgumentNotNull;
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
-import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState;
 
 import java.util.Iterator;
@@ -31,6 +31,7 @@ import org.apache.beam.sdk.values.KV;
 import org.apache.beam.sdk.values.PCollection;
 import org.apache.beam.sdk.values.TimestampedValue;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
+import org.checkerframework.checker.nullness.qual.Nullable;
 
 /**
  * {@link PTransform} and {@link Combine.CombineFn} for computing the latest 
element in a {@link
@@ -50,9 +51,6 @@ import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.Vi
  *
  * <p>For elements with the same timestamp, the element chosen for output is 
arbitrary.
  */
-@SuppressWarnings({
-  "nullness" // TODO(https://github.com/apache/beam/issues/20497)
-})
 public class Latest {
   // Do not instantiate
   private Latest() {}
@@ -88,25 +86,37 @@ public class Latest {
   /**
    * A {@link Combine.CombineFn} that computes the latest element from a set 
of inputs.
    *
+   * <p>The accumulator holds a {@literal null} value until the first input 
arrives, so combining an
+   * empty input yields {@literal null}. {@code T} must therefore be 
instantiated at a nullable type
+   * for the output to be sound, which the declaration cannot say without 
changing the public
+   * signature of {@link Latest#combineFn()}.
+   *
    * @param <T> Type of input element.
    * @see Latest
    */
   @VisibleForTesting
-  static class LatestFn<T> extends Combine.CombineFn<TimestampedValue<T>, 
TimestampedValue<T>, T> {
+  static class LatestFn<T>
+      extends Combine.CombineFn<TimestampedValue<T>, 
TimestampedValue<@Nullable T>, T> {
     /** Construct a new {@link LatestFn} instance. */
     public LatestFn() {}
 
     @Override
-    public TimestampedValue<T> createAccumulator() {
+    public TimestampedValue<@Nullable T> createAccumulator() {
       return TimestampedValue.atMinimumTimestamp(null);
     }
 
     @Override
-    public TimestampedValue<T> addInput(
-        TimestampedValue<T> accumulator, TimestampedValue<T> input) {
-      checkNotNull(accumulator, "accumulator must be non-null");
-      checkNotNull(input, "input must be non-null");
+    public TimestampedValue<@Nullable T> addInput(
+        TimestampedValue<@Nullable T> accumulator, TimestampedValue<T> input) {
+      checkArgumentNotNull(accumulator, "accumulator must be non-null");
+      checkArgumentNotNull(input, "input must be non-null");
+
+      return latest(accumulator, input);
+    }
 
+    /** Returns whichever argument is later, preferring a non-null value when 
they are equal. */
+    private static <T> TimestampedValue<@Nullable T> latest(
+        TimestampedValue<@Nullable T> accumulator, TimestampedValue<@Nullable 
T> input) {
       if (input.getTimestamp().isBefore(accumulator.getTimestamp())) {
         return accumulator;
       } else if (input.getTimestamp().isAfter(accumulator.getTimestamp())) {
@@ -117,13 +127,15 @@ public class Latest {
     }
 
     @Override
-    public Coder<TimestampedValue<T>> getAccumulatorCoder(
+    @SuppressWarnings("nullness") // accumulated values may be null
+    public Coder<TimestampedValue<@Nullable T>> getAccumulatorCoder(
         CoderRegistry registry, Coder<TimestampedValue<T>> inputCoder)
         throws CannotProvideCoderException {
       return NullableCoder.of(inputCoder);
     }
 
     @Override
+    @SuppressWarnings("nullness") // the output is null when the input is empty
     public Coder<T> getDefaultOutputCoder(
         CoderRegistry registry, Coder<TimestampedValue<T>> inputCoder)
         throws CannotProvideCoderException {
@@ -138,24 +150,28 @@ public class Latest {
     }
 
     @Override
-    public TimestampedValue<T> mergeAccumulators(Iterable<TimestampedValue<T>> 
accumulators) {
-      checkNotNull(accumulators, "accumulators must be non-null");
+    public TimestampedValue<@Nullable T> mergeAccumulators(
+        Iterable<TimestampedValue<@Nullable T>> accumulators) {
+      checkArgumentNotNull(accumulators, "accumulators must be non-null");
 
-      Iterator<TimestampedValue<T>> iter = accumulators.iterator();
+      Iterator<TimestampedValue<@Nullable T>> iter = accumulators.iterator();
       if (!iter.hasNext()) {
         return createAccumulator();
       }
 
-      TimestampedValue<T> merged = iter.next();
+      TimestampedValue<@Nullable T> merged = iter.next();
       while (iter.hasNext()) {
-        merged = addInput(merged, iter.next());
+        TimestampedValue<@Nullable T> next = iter.next();
+        checkArgumentNotNull(next, "input must be non-null");
+        merged = latest(merged, next);
       }
 
       return merged;
     }
 
     @Override
-    public T extractOutput(TimestampedValue<T> accumulator) {
+    @SuppressWarnings("nullness") // the output is null until the first input 
arrives
+    public T extractOutput(TimestampedValue<@Nullable T> accumulator) {
       return accumulator.getValue();
     }
   }
@@ -178,7 +194,7 @@ public class Latest {
       extends PTransform<PCollection<KV<K, V>>, PCollection<KV<K, V>>> {
     @Override
     public PCollection<KV<K, V>> expand(PCollection<KV<K, V>> input) {
-      checkNotNull(input);
+      checkArgumentNotNull(input);
       checkArgument(
           input.getCoder() instanceof KvCoder,
           "Input specifiedCoder must be an instance of KvCoder, but was %s",
diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java
index da6feef92d6..2c4f16f1227 100644
--- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java
+++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java
@@ -32,6 +32,7 @@ import org.apache.beam.sdk.values.TimestampedValue;
 import org.apache.beam.sdk.values.TimestampedValue.TimestampedValueCoder;
 import org.apache.beam.sdk.values.ValueInSingleWindow;
 import org.apache.beam.sdk.values.ValueKind;
+import org.checkerframework.checker.nullness.qual.Nullable;
 import org.joda.time.Duration;
 import org.joda.time.Instant;
 
@@ -39,9 +40,6 @@ import org.joda.time.Instant;
  * {@link PTransform PTransforms} for converting between explicit and implicit 
form of various Beam
  * values.
  */
-@SuppressWarnings({
-  "nullness" // TODO(https://github.com/apache/beam/issues/20497)
-})
 public class Reify {
   private static class ReifyView<K, V> extends PTransform<PCollection<K>, 
PCollection<KV<K, V>>> {
     private final PCollectionView<V> view;
@@ -80,7 +78,7 @@ public class Reify {
     @Override
     public PCollection<V> expand(PBegin input) {
       return input
-          .apply(Create.of((Void) null).withCoder(VoidCoder.of()))
+          .apply(Create.<@Nullable Void>of((@Nullable Void) 
null).withCoder(VoidCoder.of()))
           .apply(Reify.viewAsValues(view, coder))
           .apply(Values.create());
     }
diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Sample.java 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Sample.java
index d0017c88f6a..a838106ff37 100644
--- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Sample.java
+++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Sample.java
@@ -43,10 +43,7 @@ import org.apache.beam.sdk.values.PCollection;
  * <p>{@link #combineFn} can also be used manually, in combination with state 
and with the {@link
  * Combine} transform.
  */
-@SuppressWarnings({
-  "nullness", // TODO(https://github.com/apache/beam/issues/20497)
-  "rawtypes"
-})
+@SuppressWarnings({"rawtypes"})
 public class Sample {
 
   /** Returns a {@link CombineFn} that computes a fixed-sized uniform sample 
of its inputs. */
@@ -282,6 +279,7 @@ public class Sample {
     }
 
     @Override
+    @SuppressWarnings("nullness") // the output is null when the input is empty
     public T extractOutput(List<T> accumulator) {
       Iterator<T> it = internal.extractOutput(accumulator).iterator();
       return it.hasNext() ? it.next() : null;
diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/TimestampedValue.java 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/TimestampedValue.java
index 753ee4c2b9f..8bcf9ffb0f2 100644
--- 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/TimestampedValue.java
+++ 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/TimestampedValue.java
@@ -31,6 +31,7 @@ import org.apache.beam.sdk.coders.InstantCoder;
 import org.apache.beam.sdk.coders.StructuredCoder;
 import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
 import org.checkerframework.checker.nullness.qual.Nullable;
+import org.checkerframework.framework.qual.Covariant;
 import org.joda.time.Instant;
 
 /**
@@ -41,6 +42,8 @@ import org.joda.time.Instant;
  *
  * @param <V> the type of the value
  */
+// Immutable, so the value type may be widened.
+@Covariant(0)
 public class TimestampedValue<V extends @Nullable Object> {
   /**
    * Returns a new {@link TimestampedValue} with the {@link 
BoundedWindow#TIMESTAMP_MIN_VALUE
diff --git 
a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/LatestFnTest.java 
b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/LatestFnTest.java
index 3eafbaa7c86..b37f0588438 100644
--- 
a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/LatestFnTest.java
+++ 
b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/LatestFnTest.java
@@ -100,14 +100,14 @@ public class LatestFnTest {
 
   @Test
   public void testAddInputNullAccumulator() {
-    thrown.expect(NullPointerException.class);
+    thrown.expect(IllegalArgumentException.class);
     thrown.expectMessage("accumulator");
     fn.addInput(null, TV);
   }
 
   @Test
   public void testAddInputNullInput() {
-    thrown.expect(NullPointerException.class);
+    thrown.expect(IllegalArgumentException.class);
     thrown.expectMessage("input");
     fn.addInput(TV, null);
   }
@@ -151,7 +151,7 @@ public class LatestFnTest {
 
   @Test
   public void testMergeAccumulatorsNullIterable() {
-    thrown.expect(NullPointerException.class);
+    thrown.expect(IllegalArgumentException.class);
     thrown.expectMessage("accumulators");
     fn.mergeAccumulators(null);
   }

Reply via email to