1717 */
1818package org .apache .beam .sdk .extensions .avro .coders ;
1919
20+ import static org .apache .beam .vendor .guava .v32_1_2_jre .com .google .common .base .Preconditions .checkNotNull ;
21+
2022import com .google .errorprone .annotations .FormatMethod ;
2123import com .google .errorprone .annotations .FormatString ;
2224import edu .umd .cs .findbugs .annotations .SuppressFBWarnings ;
109111 *
110112 * @param <T> the type of elements handled by this coder
111113 */
112- @ SuppressWarnings ({
113- "nullness" // TODO(https://github.com/apache/beam/issues/20497)
114- })
115114public class AvroCoder <T > extends CustomCoder <T > {
116115
117116 private static final Cache <AvroCoderCacheKey , AvroCoder <?>> AVRO_CODER_CACHE =
@@ -137,7 +136,12 @@ public static <T> AvroCoder<T> specific(TypeDescriptor<T> type) {
137136 * suite for encoding and decoding.
138137 */
139138 public static <T > AvroCoder <T > specific (Class <T > type ) {
140- return specific (type , new SpecificData (type .getClassLoader ()).getSchema (type ));
139+ return specific (type , specificSchemaOf (type ));
140+ }
141+
142+ @ SuppressWarnings ("nullness" ) // SpecificData tolerates a null class loader but is unannotated
143+ private static Schema specificSchemaOf (Class <?> type ) {
144+ return new SpecificData (type .getClassLoader ()).getSchema (type );
141145 }
142146
143147 /**
@@ -167,7 +171,12 @@ public static <T> AvroCoder<T> reflect(TypeDescriptor<T> type) {
167171 * suite for encoding and decoding.
168172 */
169173 public static <T > AvroCoder <T > reflect (Class <T > type ) {
170- return reflect (type , new ReflectData (type .getClassLoader ()).getSchema (type ));
174+ return reflect (type , reflectSchemaOf (type ));
175+ }
176+
177+ @ SuppressWarnings ("nullness" ) // ReflectData tolerates a null class loader but is unannotated
178+ private static Schema reflectSchemaOf (Class <?> type ) {
179+ return new ReflectData (type .getClassLoader ()).getSchema (type );
171180 }
172181
173182 /**
@@ -395,10 +404,10 @@ public Schema get() {
395404
396405 // writer and reader are unused but kept for serialization update compatibility.
397406 @ SuppressWarnings ("unused" )
398- private final EmptyOnDeserializationThreadLocal <DatumWriter <T >> writer = null ;
407+ private final @ Nullable EmptyOnDeserializationThreadLocal <DatumWriter <T >> writer = null ;
399408
400409 @ SuppressWarnings ("unused" )
401- private final EmptyOnDeserializationThreadLocal <DatumReader <T >> reader = null ;
410+ private final @ Nullable EmptyOnDeserializationThreadLocal <DatumReader <T >> reader = null ;
402411
403412 // datumReader and datumWriter are initialized in the constructor and
404413 // on deserialization (see readObject).
@@ -424,7 +433,8 @@ protected AvroCoder(AvroDatumFactory<T> datumFactory, Schema schema) {
424433 this .decoder = new EmptyOnDeserializationThreadLocal <>();
425434 this .encoder = new EmptyOnDeserializationThreadLocal <>();
426435
427- initializeAvroDatumReaderAndWriter ();
436+ this .datumReader = datumFactory .apply (schema , schema );
437+ this .datumWriter = datumFactory .apply (schema );
428438 }
429439
430440 /** Returns the type this coder encodes/decodes. */
@@ -473,6 +483,11 @@ public T decode(InputStream inStream) throws IOException {
473483 BinaryDecoder decoderInstance = DECODER_FACTORY .directBinaryDecoder (inStream , decoder .get ());
474484 // Save the potentially-new instance for later.
475485 decoder .set (decoderInstance );
486+ return readWithoutReuse (decoderInstance );
487+ }
488+
489+ @ SuppressWarnings ("nullness" ) // DatumReader.read accepts a null reuse but is unannotated
490+ private T readWithoutReuse (BinaryDecoder decoderInstance ) throws IOException {
476491 return datumReader .read (null , decoderInstance );
477492 }
478493
@@ -808,10 +823,10 @@ private void checkMap(String context, TypeDescriptor<?> type, Schema schema) {
808823 }
809824
810825 private void checkArray (String context , TypeDescriptor <?> type , Schema schema ) {
811- TypeDescriptor <?> elementType = null ;
826+ TypeDescriptor <?> elementType ;
812827 if (type .isArray ()) {
813828 // The type is an array (with ordering)-> deterministic iff the element is deterministic.
814- elementType = type .getComponentType ();
829+ elementType = checkNotNull ( type .getComponentType () );
815830 } else if (isSubtypeOf (type , Collection .class )) {
816831 if (isSubtypeOf (type , List .class , SortedSet .class )) {
817832 // Ordered collection -> deterministic iff the element is deterministic
@@ -895,16 +910,12 @@ private void readObject(ObjectInputStream in) throws IOException, ClassNotFoundE
895910 this .datumReader = cachedCoder .get ().datumReader ;
896911 this .datumWriter = cachedCoder .get ().datumWriter ;
897912 } else {
898- initializeAvroDatumReaderAndWriter ();
913+ Schema schema = this .schemaSupplier .get ();
914+ this .datumReader = this .datumFactory .apply (schema , schema );
915+ this .datumWriter = this .datumFactory .apply (schema );
899916 }
900917 }
901918
902- private void initializeAvroDatumReaderAndWriter () {
903- this .datumReader =
904- this .datumFactory .apply (this .schemaSupplier .get (), this .schemaSupplier .get ());
905- this .datumWriter = this .datumFactory .apply (this .schemaSupplier .get ());
906- }
907-
908919 enum AvroCoderType {
909920 SPECIFIC ,
910921 REFLECT ;
0 commit comments