From ed30e2efce9141a540bdb36e256355fb9c834914 Mon Sep 17 00:00:00 2001 From: Yuan-Fang Lin Date: Fri, 2 Oct 2026 13:43:17 -0700 Subject: [PATCH] Clear all 48 oss-canary CVEs blocking ELR and bump to 2.45.61 ELR 367778 for beam-runners-flink-1.18:2.45.60 failed oss-canary validation with 48 flagged Maven coordinates. This clears all 48 and shrinks the runtime dependency graph from 212 to 70 nodes. Drop google-cloud-platform-core from the Flink runner (clears 31). It was a compile-scope dependency reaching the published POM, but its only use in main source was one line in FlinkJobServerDriver setting a GCS upload buffer. beam-runner already excludes the GCP artifacts. Bump flagged libraries: jackson 2.14.1 -> 2.19.2 (also pulls snakeyaml 2.4), classgraph -> 4.8.179, commons-compress -> 1.27.1, commons-lang3 -> 3.18.0, snappy-java -> 1.1.10.7. commons-io -> 2.16.1 and commons-codec -> 1.17.1 are required floors for commons-compress 1.26+. Exclude commons-compress from the Flink modules. flink-core declares 1.21 in its own POM and reaches it only through the compressed FileInputFormat factories, which Beam pipelines do not use. Migrate Avro 1.8.2 -> 1.11.4, which clears the last two coordinates. Move the codegen plugin to com.github.davidmc24.gradle.plugin:gradle-avro-plugin:1.9.1 so generated sources match the runtime Avro version, declare org.tukaani:xz explicitly since Avro 1.11 marks it optional, and add java.time branches to the AvroUtils converters for Avro 1.9+ date/timestamp logical types. The Kryo and SerializableSchemaSupplier changes are ported from upstream Beam 2.75. Avro ReflectData field ordering changed in Avro 1.10, but beam-runner already pins 1.11.3, so LinkedIn runtime behavior is unchanged; the updated test expectations align Beam with what already executes. Bump version to 2.45.61. BeamModulePlugin.groovy is the publish trigger and holds project.version, so a change there without a bump would republish 2.45.60. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- buildSrc/build.gradle.kts | 7 +- .../beam/gradle/BeamModulePlugin.groovy | 25 +++--- gradle.properties | 4 +- runners/flink/flink_runner.gradle | 27 ++++--- .../runners/flink/FlinkJobServerDriver.java | 3 - sdks/java/core/build.gradle | 1 + .../org/apache/beam/sdk/coders/AvroCoder.java | 4 +- .../beam/sdk/schemas/utils/AvroUtils.java | 77 +++++++++++++++++++ .../apache/beam/sdk/coders/AvroCoderTest.java | 31 +++++--- .../beam/sdk/schemas/AvroSchemaTest.java | 75 +++++++++++------- .../beam/sdk/schemas/utils/AvroUtilsTest.java | 2 +- 11 files changed, 185 insertions(+), 71 deletions(-) diff --git a/buildSrc/build.gradle.kts b/buildSrc/build.gradle.kts index e373fb2efd09..ba85ef18a87f 100644 --- a/buildSrc/build.gradle.kts +++ b/buildSrc/build.gradle.kts @@ -30,11 +30,6 @@ repositories { url = uri("https://repo.spring.io/plugins-release/") content { includeGroup("io.spring.gradle") } } - // For obsolete Avro plugin - maven { - url = uri("https://jitpack.io") - content { includeGroup("com.github.davidmc24.gradle-avro-plugin") } - } } // Dependencies on other plugins used when this plugin is invoked @@ -45,7 +40,7 @@ dependencies { implementation("com.github.spotbugs.snom:spotbugs-gradle-plugin:5.0.3") runtimeOnly("com.google.protobuf:protobuf-gradle-plugin:0.8.13") // Enable proto code generation - runtimeOnly("com.github.davidmc24.gradle-avro-plugin:gradle-avro-plugin:0.16.0") // Enable Avro code generation + runtimeOnly("com.github.davidmc24.gradle.plugin:gradle-avro-plugin:1.9.1") // Enable Avro code generation runtimeOnly("com.diffplug.spotless:spotless-plugin-gradle:5.6.1") // Enable a code formatting plugin runtimeOnly("gradle.plugin.com.palantir.gradle.docker:gradle-docker:0.22.0") // Enable building Docker containers runtimeOnly("gradle.plugin.com.dorongold.plugins:task-tree:1.5") // Adds a 'taskTree' task to print task dependency tree diff --git a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy index 70aa916462b9..580b4ddc65c6 100644 --- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy +++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy @@ -398,7 +398,7 @@ class BeamModulePlugin implements Plugin { // Automatically use the official release version if we are performing a release // otherwise append '-SNAPSHOT' - project.version = '2.45.60' + project.version = '2.45.61' if (isLinkedin(project)) { project.ext.mavenGroupId = 'com.linkedin.beam' } @@ -473,7 +473,7 @@ class BeamModulePlugin implements Plugin { def cassandra_driver_version = "3.10.2" def cdap_version = "6.5.1" def checkerframework_version = "3.27.0" - def classgraph_version = "4.8.104" + def classgraph_version = "4.8.179" def dbcp2_version = "2.8.0" def errorprone_version = "2.10.0" // Try to keep gax_version consistent with gax-grpc version in google_cloud_platform_libraries_bom @@ -492,7 +492,7 @@ class BeamModulePlugin implements Plugin { def influxdb_version = "2.19" def httpclient_version = "4.5.13" def httpcore_version = "4.4.14" - def jackson_version = "2.14.1" + def jackson_version = "2.19.2" def jaxb_api_version = "2.3.3" def jsr305_version = "3.0.2" def everit_json_version = "1.14.1" @@ -536,8 +536,8 @@ class BeamModulePlugin implements Plugin { antlr_runtime : "org.antlr:antlr4-runtime:4.7", args4j : "args4j:args4j:2.33", auto_value_annotations : "com.google.auto.value:auto-value-annotations:$autovalue_version", - avro : "org.apache.avro:avro:1.8.2", - avro_tests : "org.apache.avro:avro:1.8.2:tests", + avro : "org.apache.avro:avro:1.11.4", + avro_tests : "org.apache.avro:avro:1.11.3:tests", aws_java_sdk_cloudwatch : "com.amazonaws:aws-java-sdk-cloudwatch:$aws_java_sdk_version", aws_java_sdk_core : "com.amazonaws:aws-java-sdk-core:$aws_java_sdk_version", aws_java_sdk_dynamodb : "com.amazonaws:aws-java-sdk-dynamodb:$aws_java_sdk_version", @@ -579,12 +579,12 @@ class BeamModulePlugin implements Plugin { cdap_plugin_zendesk : "io.cdap.plugin:zendesk-plugins:1.0.0", checker_qual : "org.checkerframework:checker-qual:$checkerframework_version", classgraph : "io.github.classgraph:classgraph:$classgraph_version", - commons_codec : "commons-codec:commons-codec:1.15", + commons_codec : "commons-codec:commons-codec:1.17.1", commons_collections : "commons-collections:commons-collections:3.2.2", - commons_compress : "org.apache.commons:commons-compress:1.21", + commons_compress : "org.apache.commons:commons-compress:1.27.1", commons_csv : "org.apache.commons:commons-csv:1.8", - commons_io : "commons-io:commons-io:2.7", - commons_lang3 : "org.apache.commons:commons-lang3:3.9", + commons_io : "commons-io:commons-io:2.16.1", + commons_lang3 : "org.apache.commons:commons-lang3:3.18.0", commons_logging : "commons-logging:commons-logging:1.2", commons_math3 : "org.apache.commons:commons-math3:3.6.1", datasketches : "org.apache.datasketches:datasketches-java:6.1.1", @@ -740,7 +740,7 @@ class BeamModulePlugin implements Plugin { slf4j_jul_to_slf4j : "org.slf4j:jul-to-slf4j:$slf4j_version", slf4j_log4j12 : "org.slf4j:slf4j-log4j12:$slf4j_version", slf4j_jcl : "org.slf4j:slf4j-jcl:$slf4j_version", - snappy_java : "org.xerial.snappy:snappy-java:1.1.8.4", + snappy_java : "org.xerial.snappy:snappy-java:1.1.10.7", spark_core : "org.apache.spark:spark-core_2.11:$spark2_version", spark_network_common : "org.apache.spark:spark-network-common_2.11:$spark2_version", spark_sql : "org.apache.spark:spark-sql_2.11:$spark2_version", @@ -764,6 +764,9 @@ class BeamModulePlugin implements Plugin { vendored_guava_26_0_jre : "org.apache.beam:beam-vendor-guava-26_0-jre:0.1", vendored_calcite_1_28_0 : "org.apache.beam:beam-vendor-calcite-1_28_0:0.2", woodstox_core_asl : "org.codehaus.woodstox:woodstox-core-asl:4.4.1", + // Avro 1.9+ marks xz optional, so it no longer arrives transitively. + // Declared explicitly to keep the xz codec working as it did under Avro 1.8. + xz : "org.tukaani:xz:1.10", zstd_jni : "com.github.luben:zstd-jni:1.5.2-5", quickcheck_core : "com.pholser:junit-quickcheck-core:$quickcheck_version", quickcheck_generators : "com.pholser:junit-quickcheck-generators:$quickcheck_version", @@ -2278,7 +2281,7 @@ class BeamModulePlugin implements Plugin { // TODO: Decide whether this should be inlined into the one project that relies on it // or be left here. - project.ext.applyAvroNature = { project.apply plugin: "com.commercehub.gradle.plugin.avro" } + project.ext.applyAvroNature = { project.apply plugin: "com.github.davidmc24.gradle.plugin.avro" } project.ext.applyAntlrNature = { project.apply plugin: 'antlr' diff --git a/gradle.properties b/gradle.properties index b8863f56a2fc..142ebc670c02 100644 --- a/gradle.properties +++ b/gradle.properties @@ -30,8 +30,8 @@ signing.gnupg.useLegacyGpg=true # buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy. # To build a custom Beam version make sure you change it in both places, see # https://github.com/apache/beam/issues/21302. -version=2.45.60 -sdk_version=2.45.60 +version=2.45.61 +sdk_version=2.45.61 javaVersion=1.8 diff --git a/runners/flink/flink_runner.gradle b/runners/flink/flink_runner.gradle index f2bc848756f3..5cff39fbd158 100644 --- a/runners/flink/flink_runner.gradle +++ b/runners/flink/flink_runner.gradle @@ -159,24 +159,33 @@ dependencies { implementation project(":runners:core-construction-java") implementation project(":runners:java-fn-execution") implementation project(":runners:java-job-service") - implementation project(":sdks:java:extensions:google-cloud-platform-core") implementation library.java.vendored_grpc_1_48_1 implementation library.java.slf4j_api implementation library.java.joda_time implementation library.java.args4j + // flink-core declares commons-compress 1.21 (CVE-2024-25710, CVE-2024-26308). + // Flink reaches it only through the compressed FileInputFormat factories + // (Bzip2/XZ/ZStandard InputStreamFactory), which Beam pipelines never use. + // Drop the edge from every Flink module so the vulnerable coordinate leaves + // the published dependency graph; Beam's own patched commons-compress still + // provides the classes via sdks:java:core. + def excludeCommonsCompress = { + exclude group: "org.apache.commons", module: "commons-compress" + } + // flink-core-api is introduced in Flink 1.20+ if (flink_major == '1.20' || flink_major.startsWith('2')) { - implementation "org.apache.flink:flink-core-api:$flink_version" + implementation("org.apache.flink:flink-core-api:$flink_version", excludeCommonsCompress) } - implementation "org.apache.flink:flink-clients:$flink_version" + implementation("org.apache.flink:flink-clients:$flink_version", excludeCommonsCompress) // Runtime dependencies are not included in Beam's generated pom.xml, so we must declare flink-clients in implementation // configuration (https://issues.apache.org/jira/browse/BEAM-11732). permitUnusedDeclared "org.apache.flink:flink-clients:$flink_version" - implementation "org.apache.flink:flink-streaming-java:$flink_version" - implementation "org.apache.flink:flink-table:$flink_version" + implementation("org.apache.flink:flink-streaming-java:$flink_version", excludeCommonsCompress) + implementation("org.apache.flink:flink-table:$flink_version", excludeCommonsCompress) // RocksDB state backend (included in the Flink distribution) provided "org.apache.flink:flink-statebackend-rocksdb:$flink_version" testImplementation "org.apache.flink:flink-statebackend-rocksdb:$flink_version" @@ -188,13 +197,13 @@ dependencies { miniCluster "org.apache.flink:flink-runtime-web:$flink_version" - implementation "org.apache.flink:flink-core:$flink_version" + implementation("org.apache.flink:flink-core:$flink_version", excludeCommonsCompress) implementation "org.apache.flink:flink-metrics-core:$flink_version" - implementation "org.apache.flink:flink-java:$flink_version" + implementation("org.apache.flink:flink-java:$flink_version", excludeCommonsCompress) - implementation "org.apache.flink:flink-runtime:$flink_version" + implementation("org.apache.flink:flink-runtime:$flink_version", excludeCommonsCompress) implementation "org.apache.flink:flink-metrics-core:$flink_version" - implementation "org.apache.flink:flink-optimizer:$flink_version" + implementation("org.apache.flink:flink-optimizer:$flink_version", excludeCommonsCompress) testImplementation "org.apache.flink:flink-runtime:$flink_version:tests" testImplementation "org.apache.flink:flink-rpc-akka:$flink_version" testImplementation project(path: ":sdks:java:core", configuration: "shadowTest") diff --git a/runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkJobServerDriver.java b/runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkJobServerDriver.java index 671cc2597cb2..286e68575d89 100644 --- a/runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkJobServerDriver.java +++ b/runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkJobServerDriver.java @@ -18,7 +18,6 @@ package org.apache.beam.runners.flink; import org.apache.beam.runners.jobsubmission.JobServerDriver; -import org.apache.beam.sdk.extensions.gcp.options.GcsOptions; import org.apache.beam.sdk.fn.server.ServerFactory; import org.apache.beam.sdk.io.FileSystems; import org.apache.beam.sdk.options.PipelineOptions; @@ -70,8 +69,6 @@ String getFlinkConfDir() { public static void main(String[] args) throws Exception { // TODO: Expose the fileSystem related options. PipelineOptions options = PipelineOptionsFactory.create(); - // Limiting gcs upload buffer to reduce memory usage while doing parallel artifact uploads. - options.as(GcsOptions.class).setGcsUploadBufferSizeBytes(1024 * 1024); // Register standard file systems. FileSystems.setDefaultPipelineOptions(options); fromParams(args).run(); diff --git a/sdks/java/core/build.gradle b/sdks/java/core/build.gradle index 2e172ec50c07..6cc7711c653d 100644 --- a/sdks/java/core/build.gradle +++ b/sdks/java/core/build.gradle @@ -90,6 +90,7 @@ dependencies { shadow library.java.jackson_databind shadow library.java.slf4j_api shadow library.java.avro + shadow library.java.xz shadow library.java.snappy_java shadow library.java.joda_time implementation enforcedPlatform(library.java.google_cloud_platform_libraries_bom) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/coders/AvroCoder.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/coders/AvroCoder.java index c7b39d5b025a..c5ec89ffd70c 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/coders/AvroCoder.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/coders/AvroCoder.java @@ -251,7 +251,9 @@ private Object readResolve() throws IOException, ClassNotFoundException { * Serializable}'s usage of the {@link #writeReplace} method. Kryo doesn't utilize Java's * serialization and hence is able to encode the {@link Schema} object directly. */ - private static class SerializableSchemaSupplier implements Serializable, Supplier { + // Package-private (not private) so tests can register a Kryo JavaSerializer for it. + // Matches upstream Beam 2.75. + static class SerializableSchemaSupplier implements Serializable, Supplier { // writeReplace makes this object serializable. This is a limitation of FindBugs as discussed // here: // http://stackoverflow.com/questions/26156523/is-writeobject-not-neccesary-using-the-serialization-proxy-pattern diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/utils/AvroUtils.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/utils/AvroUtils.java index be387425f511..e9a3f2162e68 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/utils/AvroUtils.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/utils/AvroUtils.java @@ -238,6 +238,68 @@ public org.apache.avro.Schema toAvroType(String name, String namespace) { } } + /** + * Conversions between the {@code java.time} types that Avro 1.9+ generates for date and timestamp + * logical types and the Joda types that Beam schemas use. + * + *

These are invoked from generated bytecode in {@link AvroConvertValueForGetter} and {@link + * AvroConvertValueForSetter}. Avro 1.8 generated Joda types directly, so no conversion was + * needed; from 1.9 onward Joda support was removed in favour of {@code java.time}. + */ + public static @Nullable Instant javaLocalDateToJoda(@Nullable Object value) { + if (value == null) { + return null; + } + return new Instant( + ((java.time.LocalDate) value) + .atStartOfDay(java.time.ZoneOffset.UTC) + .toInstant() + .toEpochMilli()); + } + + public static @Nullable Instant javaInstantToJoda(@Nullable Object value) { + if (value == null) { + return null; + } + return new Instant(((java.time.Instant) value).toEpochMilli()); + } + + public static java.time.@Nullable LocalDate jodaToJavaLocalDate(@Nullable Object value) { + if (value == null) { + return null; + } + return java.time.Instant.ofEpochMilli(((ReadableInstant) value).getMillis()) + .atZone(java.time.ZoneOffset.UTC) + .toLocalDate(); + } + + public static java.time.@Nullable Instant jodaToJavaInstant(@Nullable Object value) { + if (value == null) { + return null; + } + return java.time.Instant.ofEpochMilli(((ReadableInstant) value).getMillis()); + } + + /** Generates a call to one of the static conversion helpers above. */ + private static StackManipulation invokeAvroTimeConversion( + StackManipulation readValue, String methodName) { + return new Compound( + readValue, + MethodInvocation.invoke( + new ForLoadedType(AvroUtils.class) + .getDeclaredMethods() + .filter(ElementMatchers.named(methodName).and(ElementMatchers.isStatic())) + .getOnly())); + } + + private static boolean isJavaTimeLocalDate(TypeDescriptor type) { + return java.time.LocalDate.class.equals(type.getRawType()); + } + + private static boolean isJavaTimeInstant(TypeDescriptor type) { + return java.time.Instant.class.equals(type.getRawType()); + } + public static class AvroConvertType extends ConvertType { public AvroConvertType(boolean returnRawType) { super(returnRawType); @@ -247,6 +309,10 @@ public AvroConvertType(boolean returnRawType) { protected java.lang.reflect.Type convertDefault(TypeDescriptor type) { if (type.isSubtypeOf(TypeDescriptor.of(GenericFixed.class))) { return byte[].class; + } else if (isJavaTimeLocalDate(type) || isJavaTimeInstant(type)) { + // Avro 1.9+ date/timestamp logical types surface as java.time; Beam models them as + // Joda Instant. + return Instant.class; } else { return super.convertDefault(type); } @@ -277,6 +343,10 @@ protected StackManipulation convertDefault(TypeDescriptor type) { ElementMatchers.named("bytes") .and(ElementMatchers.returns(new ForLoadedType(byte[].class)))) .getOnly())); + } else if (isJavaTimeLocalDate(type)) { + return invokeAvroTimeConversion(readValue, "javaLocalDateToJoda"); + } else if (isJavaTimeInstant(type)) { + return invokeAvroTimeConversion(readValue, "javaInstantToJoda"); } return super.convertDefault(type); } @@ -313,6 +383,10 @@ protected StackManipulation convertDefault(TypeDescriptor type) { ElementMatchers.isConstructor() .and(ElementMatchers.takesArguments(byteArrayType))) .getOnly())); + } else if (isJavaTimeLocalDate(type)) { + return invokeAvroTimeConversion(readValue, "jodaToJavaLocalDate"); + } else if (isJavaTimeInstant(type)) { + return invokeAvroTimeConversion(readValue, "jodaToJavaInstant"); } return super.convertDefault(type); } @@ -1186,6 +1260,9 @@ private static org.apache.avro.Schema getFieldSchema( } else if (logicalType instanceof LogicalTypes.TimestampMillis) { if (value instanceof ReadableInstant) { return convertDateTimeStrict(((ReadableInstant) value).getMillis(), fieldType); + } else if (value instanceof java.time.Instant) { + // Avro 1.9+ represents timestamp logical types with java.time. + return convertDateTimeStrict(((java.time.Instant) value).toEpochMilli(), fieldType); } else { return convertDateTimeStrict((Long) value, fieldType); } diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/coders/AvroCoderTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/coders/AvroCoderTest.java index 76e1568dfab9..3287db980c0b 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/coders/AvroCoderTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/coders/AvroCoderTest.java @@ -28,6 +28,7 @@ import com.esotericsoftware.kryo.Kryo; import com.esotericsoftware.kryo.io.Input; import com.esotericsoftware.kryo.io.Output; +import com.esotericsoftware.kryo.serializers.JavaSerializer; import java.io.ByteArrayInputStream; import java.io.ByteArrayOutputStream; import java.io.ObjectInputStream; @@ -85,7 +86,6 @@ import org.hamcrest.TypeSafeMatcher; import org.joda.time.DateTime; import org.joda.time.DateTimeZone; -import org.joda.time.LocalDate; import org.junit.Rule; import org.junit.Test; import org.junit.experimental.categories.Category; @@ -112,8 +112,9 @@ public class AvroCoderTest { "mystring", ByteBuffer.wrap(new byte[] {1, 2, 3, 4}), new fixed4(new byte[] {1, 2, 3, 4}), - new LocalDate(1979, 3, 14), - new DateTime().withDate(1979, 3, 14).withTime(1, 2, 3, 4), + java.time.LocalDate.of(1979, 3, 14), + java.time.Instant.ofEpochMilli( + new DateTime().withDate(1979, 3, 14).withTime(1, 2, 3, 4).getMillis()), TestEnum.abc, AVRO_NESTED_SPECIFIC_RECORD, ImmutableList.of(AVRO_NESTED_SPECIFIC_RECORD, AVRO_NESTED_SPECIFIC_RECORD), @@ -283,12 +284,17 @@ public void testKryoSerialization() throws Exception { // Kryo instantiation Kryo kryo = new Kryo(); + kryo.setRegistrationRequired(false); kryo.setInstantiatorStrategy(new StdInstantiatorStrategy()); + // Avro 1.9+ gives JsonProperties an instance-level immutable `reserved` Set that Kryo + // cannot rebuild. Delegate the schema supplier to Java serialization, which is what it + // is designed for. Matches upstream Beam 2.75. + kryo.addDefaultSerializer(AvroCoder.SerializableSchemaSupplier.class, JavaSerializer.class); // Serialization of object without any memoization ByteArrayOutputStream coderWithoutMemoizationBos = new ByteArrayOutputStream(); try (Output output = new Output(coderWithoutMemoizationBos)) { - kryo.writeObject(output, coder); + kryo.writeClassAndObject(output, coder); } // Force thread local memoization to store values. @@ -297,18 +303,18 @@ public void testKryoSerialization() throws Exception { // Serialization of object with memoized fields ByteArrayOutputStream coderWithMemoizationBos = new ByteArrayOutputStream(); try (Output output = new Output(coderWithMemoizationBos)) { - kryo.writeObject(output, coder); + kryo.writeClassAndObject(output, coder); } // Copy empty and memoized variants of the Coder ByteArrayInputStream bisWithoutMemoization = new ByteArrayInputStream(coderWithoutMemoizationBos.toByteArray()); AvroCoder copiedWithoutMemoization = - (AvroCoder) kryo.readObject(new Input(bisWithoutMemoization), AvroCoder.class); + (AvroCoder) kryo.readClassAndObject(new Input(bisWithoutMemoization)); ByteArrayInputStream bisWithMemoization = new ByteArrayInputStream(coderWithMemoizationBos.toByteArray()); AvroCoder copiedWithMemoization = - (AvroCoder) kryo.readObject(new Input(bisWithMemoization), AvroCoder.class); + (AvroCoder) kryo.readClassAndObject(new Input(bisWithMemoization)); CoderProperties.coderDecodeEncodeEqual(copiedWithoutMemoization, value); CoderProperties.coderDecodeEncodeEqual(copiedWithMemoization, value); @@ -350,11 +356,12 @@ public void testDisableReflectionEncoding() { AvroCoder.of(Pojo.class, false); fail("When userReclectApi is disable, schema should not be generated through reflection"); } catch (AvroRuntimeException e) { - String message = - "avro.shaded.com.google.common.util.concurrent.UncheckedExecutionException: " - + "org.apache.avro.AvroRuntimeException: " - + "Not a Specific class: class org.apache.beam.sdk.coders.AvroCoderTest$Pojo"; - assertEquals(message, e.getMessage()); + // Avro 1.9+ dropped its shaded Guava, so the exception is no longer wrapped in + // avro.shaded...UncheckedExecutionException. Assert on the meaningful part only. + assertThat( + e.getMessage(), + containsString( + "Not a Specific class: class org.apache.beam.sdk.coders.AvroCoderTest$Pojo")); } } diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/AvroSchemaTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/AvroSchemaTest.java index 417d45b48f52..802b3dde78ad 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/AvroSchemaTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/AvroSchemaTest.java @@ -259,6 +259,15 @@ public String toString() { .build(); private static final FieldType SUB_TYPE = FieldType.row(SUBSCHEMA).withNullable(true); + // The POJO path derives its schema by reflection, which Avro 1.9+ orders by Java field name + // (anInt, boolNonNullable). The SpecificRecord path above keeps the .avsc declaration order. + private static final Schema POJO_SUBSCHEMA = + Schema.builder() + .addNullableField("int", FieldType.INT32) + .addField("BOOL_NON_NULLABLE", FieldType.BOOLEAN) + .build(); + private static final FieldType POJO_SUB_TYPE = FieldType.row(POJO_SUBSCHEMA).withNullable(true); + private static final EnumerationType TEST_ENUM_TYPE = EnumerationType.create("abc", "cde"); private static final Schema SCHEMA = @@ -279,28 +288,39 @@ public String toString() { .addNullableField("map", FieldType.map(FieldType.STRING, SUB_TYPE)) .build(); + // Avro 1.9+ derives reflect schemas with a deterministic field order (sorted by Java field + // name) instead of the JVM-dependent declaration order used by Avro 1.8. The field set is + // unchanged; only the order differs. Note this shifts Beam Row encoding positions for + // POJO-derived Avro schemas. private static final Schema POJO_SCHEMA = Schema.builder() - .addField("bool_non_nullable", FieldType.BOOLEAN) - .addNullableField("int", FieldType.INT32) - .addNullableField("long", FieldType.INT64) - .addNullableField("float", FieldType.FLOAT) .addNullableField("double", FieldType.DOUBLE) - .addNullableField("string", FieldType.STRING) + .addNullableField("float", FieldType.FLOAT) + .addNullableField("long", FieldType.INT64) + .addNullableField("int", FieldType.INT32) + .addNullableField("array", FieldType.array(POJO_SUB_TYPE.withNullable(false))) + .addField("bool_non_nullable", FieldType.BOOLEAN) .addNullableField("bytes", FieldType.BYTES) - .addField("fixed", FieldType.logicalType(FixedBytes.of(4))) .addField("date", FieldType.DATETIME) - .addField("timestampMillis", FieldType.DATETIME) + .addField("fixed", FieldType.logicalType(FixedBytes.of(4))) + .addNullableField( + "map", FieldType.map(FieldType.STRING, POJO_SUB_TYPE.withNullable(false))) + .addNullableField("row", POJO_SUB_TYPE) + .addNullableField("string", FieldType.STRING) .addField("testEnum", FieldType.logicalType(TEST_ENUM_TYPE)) - .addNullableField("row", SUB_TYPE) - .addNullableField("array", FieldType.array(SUB_TYPE.withNullable(false))) - .addNullableField("map", FieldType.map(FieldType.STRING, SUB_TYPE.withNullable(false))) + .addField("timestampMillis", FieldType.DATETIME) .build(); private static final byte[] BYTE_ARRAY = new byte[] {1, 2, 3, 4}; private static final DateTime DATE_TIME = new DateTime().withDate(1979, 3, 14).withTime(1, 2, 3, 4); private static final LocalDate DATE = new LocalDate(1979, 3, 14); + // Avro 1.9+ generates java.time types for date/timestamp logical types, so the + // SpecificRecord constructor needs java.time values. The Joda constants above + // remain the expected values on the Beam-schema side of each assertion. + private static final java.time.LocalDate AVRO_DATE = java.time.LocalDate.of(1979, 3, 14); + private static final java.time.Instant AVRO_DATE_TIME = + java.time.Instant.ofEpochMilli(DATE_TIME.getMillis()); private static final TestAvroNested AVRO_NESTED_SPECIFIC_RECORD = new TestAvroNested(true, 42); private static final TestAvro AVRO_SPECIFIC_RECORD = new TestAvro( @@ -312,8 +332,8 @@ public String toString() { "mystring", ByteBuffer.wrap(BYTE_ARRAY), new fixed4(BYTE_ARRAY), - DATE, - DATE_TIME, + AVRO_DATE, + AVRO_DATE_TIME, TestEnum.abc, AVRO_NESTED_SPECIFIC_RECORD, ImmutableList.of(AVRO_NESTED_SPECIFIC_RECORD, AVRO_NESTED_SPECIFIC_RECORD), @@ -350,6 +370,8 @@ public String toString() { .build(); private static final Row NESTED_ROW = Row.withSchema(SUBSCHEMA).addValues(true, 42).build(); + private static final Row POJO_NESTED_ROW = + Row.withSchema(POJO_SUBSCHEMA).addValues(42, true).build(); private static final Row ROW = Row.withSchema(SCHEMA) .addValues( @@ -425,23 +447,24 @@ public void testRowToGenericRecord() { ImmutableList.of(SUB_POJO, SUB_POJO), ImmutableMap.of("k1", SUB_POJO, "k2", SUB_POJO)); + // Values must track POJO_SCHEMA's field order (sorted by Java field name under Avro 1.9+). private static final Row ROW_FOR_POJO = Row.withSchema(POJO_SCHEMA) .addValues( - true, - 43, - 44L, - (float) 44.1, - (double) 44.2, - "mystring", - ByteBuffer.wrap(BYTE_ARRAY), - BYTE_ARRAY, - DATE.toDateTimeAtStartOfDay(DateTimeZone.UTC), - DATE_TIME, - TEST_ENUM_TYPE.valueOf("abc"), - NESTED_ROW, - ImmutableList.of(NESTED_ROW, NESTED_ROW), - ImmutableMap.of("k1", NESTED_ROW, "k2", NESTED_ROW)) + (double) 44.2, // double + (float) 44.1, // float + 44L, // long + 43, // int + ImmutableList.of(POJO_NESTED_ROW, POJO_NESTED_ROW), // array + true, // bool_non_nullable + ByteBuffer.wrap(BYTE_ARRAY), // bytes + DATE.toDateTimeAtStartOfDay(DateTimeZone.UTC), // date + BYTE_ARRAY, // fixed + ImmutableMap.of("k1", POJO_NESTED_ROW, "k2", POJO_NESTED_ROW), // map + POJO_NESTED_ROW, // row + "mystring", // string + TEST_ENUM_TYPE.valueOf("abc"), // testEnum + DATE_TIME) // timestampMillis .build(); @Test diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/utils/AvroUtilsTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/utils/AvroUtilsTest.java index 4616fd487ebd..72850c02b6bf 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/utils/AvroUtilsTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/utils/AvroUtilsTest.java @@ -32,11 +32,11 @@ import org.apache.avro.Conversions; import org.apache.avro.LogicalType; import org.apache.avro.LogicalTypes; -import org.apache.avro.RandomData; import org.apache.avro.Schema.Type; import org.apache.avro.generic.GenericRecord; import org.apache.avro.generic.GenericRecordBuilder; import org.apache.avro.reflect.ReflectData; +import org.apache.avro.util.RandomData; import org.apache.avro.util.Utf8; import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.coders.AvroCoder;