Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 1 addition & 6 deletions buildSrc/build.gradle.kts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -398,7 +398,7 @@ class BeamModulePlugin implements Plugin<Project> {

// 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'
}
Expand Down Expand Up @@ -473,7 +473,7 @@ class BeamModulePlugin implements Plugin<Project> {
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
Expand All @@ -492,7 +492,7 @@ class BeamModulePlugin implements Plugin<Project> {
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"
Expand Down Expand Up @@ -536,8 +536,8 @@ class BeamModulePlugin implements Plugin<Project> {
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",
Expand Down Expand Up @@ -579,12 +579,12 @@ class BeamModulePlugin implements Plugin<Project> {
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",
Expand Down Expand Up @@ -740,7 +740,7 @@ class BeamModulePlugin implements Plugin<Project> {
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",
Expand All @@ -764,6 +764,9 @@ class BeamModulePlugin implements Plugin<Project> {
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",
Expand Down Expand Up @@ -2278,7 +2281,7 @@ class BeamModulePlugin implements Plugin<Project> {

// 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'
Expand Down
4 changes: 2 additions & 2 deletions gradle.properties
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
27 changes: 18 additions & 9 deletions runners/flink/flink_runner.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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();
Expand Down
1 change: 1 addition & 0 deletions sdks/java/core/build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<Schema> {
// Package-private (not private) so tests can register a Kryo JavaSerializer for it.
// Matches upstream Beam 2.75.
static class SerializableSchemaSupplier implements Serializable, Supplier<Schema> {
// 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*
* <p>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);
Expand All @@ -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);
}
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -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);
}
Expand Down
Loading
Loading