From fbc5b5ea5c32deccfa62c2b60aa834d729136b12 Mon Sep 17 00:00:00 2001 From: Luke Cwik Date: Mon, 18 Jul 2016 09:08:14 -0700 Subject: [PATCH] [BEAM-465] OutgoingMessageCoder should be an AtomicCoder --- .../main/java/org/apache/beam/sdk/io/PubsubUnboundedSink.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/PubsubUnboundedSink.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/PubsubUnboundedSink.java index 78758a29a959..6f2b3ac4cc34 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/PubsubUnboundedSink.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/PubsubUnboundedSink.java @@ -20,11 +20,11 @@ import static com.google.common.base.Preconditions.checkState; +import org.apache.beam.sdk.coders.AtomicCoder; import org.apache.beam.sdk.coders.BigEndianLongCoder; import org.apache.beam.sdk.coders.ByteArrayCoder; import org.apache.beam.sdk.coders.Coder; import org.apache.beam.sdk.coders.CoderException; -import org.apache.beam.sdk.coders.CustomCoder; import org.apache.beam.sdk.coders.KvCoder; import org.apache.beam.sdk.coders.NullableCoder; import org.apache.beam.sdk.coders.StringUtf8Coder; @@ -105,7 +105,7 @@ public class PubsubUnboundedSink extends PTransform, PDone> { /** * Coder for conveying outgoing messages between internal stages. */ - private static class OutgoingMessageCoder extends CustomCoder { + private static class OutgoingMessageCoder extends AtomicCoder { private static final NullableCoder RECORD_ID_CODER = NullableCoder.of(StringUtf8Coder.of());