This is an automated email from the ASF dual-hosted git repository.
Abacn 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 f3e12fee1b2 [DebeziumIO] Upgrade to Debezium 3.5.2.Final (#39569)
f3e12fee1b2 is described below
commit f3e12fee1b2beffe22ed41d82e67f3069f0a093f
Author: Tobias Kaymak <[email protected]>
AuthorDate: Mon Aug 3 21:45:26 2026 +0200
[DebeziumIO] Upgrade to Debezium 3.5.2.Final (#39569)
* Remove dead snapshot.mode override in DebeziumIO.getRecordSchema
* Upgrade DebeziumIO to Debezium 3.5.2.Final
* Wire DebeziumSDFDatabaseHistory default to schema.history.internal
* Update sdks/java/io/debezium/src/README.md
---
.../test/integration/io/xlang/debezium/debezium.go | 2 +-
.../integration/io/xlang/debezium/debezium_test.go | 2 +-
sdks/java/io/debezium/build.gradle | 17 ++++++++-----
.../io/debezium/expansion-service/build.gradle | 6 +++--
sdks/java/io/debezium/src/README.md | 8 +++---
.../org/apache/beam/io/debezium/DebeziumIO.java | 15 ++++++-----
.../beam/io/debezium/KafkaSourceConsumerFn.java | 5 ++++
.../io/debezium/DebeziumIOMySqlConnectorIT.java | 6 ++---
.../debezium/DebeziumIOPostgresSqlConnectorIT.java | 4 +--
.../apache/beam/io/debezium/DebeziumIOTest.java | 3 ++-
.../debezium/DebeziumReadSchemaTransformTest.java | 29 ++++++++++++++++------
.../io/external/xlang_debeziumio_it_test.py | 2 +-
12 files changed, 63 insertions(+), 36 deletions(-)
diff --git a/sdks/go/test/integration/io/xlang/debezium/debezium.go
b/sdks/go/test/integration/io/xlang/debezium/debezium.go
index e1b9bab963c..a4046f01874 100644
--- a/sdks/go/test/integration/io/xlang/debezium/debezium.go
+++ b/sdks/go/test/integration/io/xlang/debezium/debezium.go
@@ -31,7 +31,7 @@ func ReadPipeline(addr, username, password, dbname, host,
port string, connector
connectorClass, reflectx.String,
debeziumio.MaxRecord(maxrecords),
debeziumio.MaxTimeToRun(120000),
debeziumio.ConnectionProperties(connectionProperties),
debeziumio.ExpansionAddr(addr))
- expectedJson :=
`{"metadata":{"connector":"postgresql","version":"3.1.3.Final","name":"beam-debezium-connector","database":"inventory","schema":"inventory","table":"customers"},"before":null,"after":{"fields":{"last_name":"Thomas","id":1001,"first_name":"Sally","email":"[email protected]"}}}`
+ expectedJson :=
`{"metadata":{"connector":"postgresql","version":"3.5.2.Final","name":"beam-debezium-connector","database":"inventory","schema":"inventory","table":"customers"},"before":null,"after":{"fields":{"last_name":"Thomas","id":1001,"first_name":"Sally","email":"[email protected]"}}}`
expected := beam.Create(s, expectedJson)
passert.Equals(s, result, expected)
return p
diff --git a/sdks/go/test/integration/io/xlang/debezium/debezium_test.go
b/sdks/go/test/integration/io/xlang/debezium/debezium_test.go
index 8ccb64cae20..b234dd5d8f8 100644
--- a/sdks/go/test/integration/io/xlang/debezium/debezium_test.go
+++ b/sdks/go/test/integration/io/xlang/debezium/debezium_test.go
@@ -33,7 +33,7 @@ import (
)
const (
- debeziumImage = "quay.io/debezium/example-postgres:3.1.3.Final"
+ debeziumImage = "quay.io/debezium/example-postgres:3.5.2.Final"
debeziumPort = "5432/tcp"
maxRetries = 5
)
diff --git a/sdks/java/io/debezium/build.gradle
b/sdks/java/io/debezium/build.gradle
index c488ac17d99..ad5410b37de 100644
--- a/sdks/java/io/debezium/build.gradle
+++ b/sdks/java/io/debezium/build.gradle
@@ -46,10 +46,13 @@ dependencies {
permitUnusedDeclared library.java.jackson_dataformat_csv
// Kafka connect dependencies
- implementation "org.apache.kafka:connect-api:3.9.0"
+ implementation "org.apache.kafka:connect-api:4.1.2"
+ implementation "org.apache.kafka:kafka-clients:4.1.2"
// Debezium dependencies
- implementation group: 'io.debezium', name: 'debezium-core', version:
'3.1.3.Final'
+ implementation group: 'io.debezium', name: 'debezium-core', version:
'3.5.2.Final'
+ implementation group: 'io.debezium', name: 'debezium-config', version:
'3.5.2.Final'
+ implementation group: 'io.debezium', name: 'debezium-connector-common',
version: '3.5.2.Final'
// Test dependencies
testImplementation project(path: ":sdks:java:core", configuration:
"shadowTest")
@@ -64,12 +67,12 @@ dependencies {
testImplementation "org.testcontainers:kafka"
testImplementation "org.testcontainers:mysql"
testImplementation "org.testcontainers:postgresql"
- testImplementation
"io.debezium:debezium-testing-testcontainers:3.1.3.Final"
+ testImplementation
"io.debezium:debezium-testing-testcontainers:3.5.2.Final"
testImplementation 'com.zaxxer:HikariCP:5.1.0'
// Debezium connector implementations for testing
- testImplementation group: 'io.debezium', name: 'debezium-connector-mysql',
version: '3.1.3.Final'
- testImplementation group: 'io.debezium', name:
'debezium-connector-postgres', version: '3.1.3.Final'
+ testImplementation group: 'io.debezium', name: 'debezium-connector-mysql',
version: '3.5.2.Final'
+ testImplementation group: 'io.debezium', name:
'debezium-connector-postgres', version: '3.5.2.Final'
}
// Pin the Antlr version to 4.10
@@ -83,7 +86,9 @@ configurations.all {
'com.fasterxml.jackson.core:jackson-core:2.17.1',
'com.fasterxml.jackson.core:jackson-annotations:2.17.1',
'com.fasterxml.jackson.core:jackson-databind:2.17.1',
- 'com.fasterxml.jackson.datatype:jackson-datatype-jsr310:2.17.1'
+
'com.fasterxml.jackson.datatype:jackson-datatype-jsr310:2.17.1',
+ 'org.apache.kafka:kafka-clients:4.1.2',
+ 'org.postgresql:postgresql:42.7.7'
}
}
diff --git a/sdks/java/io/debezium/expansion-service/build.gradle
b/sdks/java/io/debezium/expansion-service/build.gradle
index 82a34b5c066..fe820b62c13 100644
--- a/sdks/java/io/debezium/expansion-service/build.gradle
+++ b/sdks/java/io/debezium/expansion-service/build.gradle
@@ -39,7 +39,7 @@ dependencies {
runtimeOnly library.java.slf4j_jdk14
// Debezium runtime dependencies
- def debezium_version = '3.1.3.Final'
+ def debezium_version = '3.5.2.Final'
runtimeOnly group: 'io.debezium', name: 'debezium-connector-mysql',
version: debezium_version
runtimeOnly group: 'io.debezium', name: 'debezium-connector-postgres',
version: debezium_version
runtimeOnly group: 'io.debezium', name: 'debezium-connector-sqlserver',
version: debezium_version
@@ -55,7 +55,9 @@ configurations.all {
'com.fasterxml.jackson.core:jackson-core:2.17.1',
'com.fasterxml.jackson.core:jackson-annotations:2.17.1',
'com.fasterxml.jackson.core:jackson-databind:2.17.1',
- 'com.fasterxml.jackson.datatype:jackson-datatype-jsr310:2.17.1'
+
'com.fasterxml.jackson.datatype:jackson-datatype-jsr310:2.17.1',
+ 'org.apache.kafka:kafka-clients:4.1.2',
+ 'org.postgresql:postgresql:42.7.7'
}
}
diff --git a/sdks/java/io/debezium/src/README.md
b/sdks/java/io/debezium/src/README.md
index 53521321885..677f42b59b0 100644
--- a/sdks/java/io/debezium/src/README.md
+++ b/sdks/java/io/debezium/src/README.md
@@ -25,7 +25,7 @@ DebeziumIO is an Apache Beam connector that lets users
connect their Events-Driv
### Getting Started
-DebeziumIO uses [Debezium Connectors
v3.1](https://debezium.io/documentation/reference/3.1/connectors/) to connect
to Apache Beam. All you need to do is choose the Debezium Connector that suits
your Debezium setup and pick a [Serializable
Function](https://beam.apache.org/releases/javadoc/2.65.0/org/apache/beam/sdk/transforms/SerializableFunction.html),
then you will be able to connect to Apache Beam and start building your own
Pipelines.
+DebeziumIO uses [Debezium Connectors
v3.5](https://debezium.io/documentation/reference/3.5/connectors/) to connect
to Apache Beam. All you need to do is choose the Debezium Connector that suits
your Debezium setup and pick a [Serializable
Function](https://beam.apache.org/releases/javadoc/2.65.0/org/apache/beam/sdk/transforms/SerializableFunction.html),
then you will be able to connect to Apache Beam and start building your own
Pipelines.
These connectors have been successfully tested and are known to work fine:
* MySQL Connector
@@ -65,7 +65,7 @@ You can also add more configuration, such as
Connector-specific Properties with
|Method|Params|Description|
|-|-|-|
|`.withConnectionProperty(propName, propValue)`|_String_, _String_|Adds a
custom property to the connector.|
-> **Note:** For more information on custom properties, see your [Debezium
Connector](https://debezium.io/documentation/reference/3.1/connectors/)
specific documentation.
+> **Note:** For more information on custom properties, see your [Debezium
Connector](https://debezium.io/documentation/reference/3.5/connectors/)
specific documentation.
Example of a MySQL Debezium Connector setup:
```
@@ -160,8 +160,8 @@ By default, DebeziumIO initializes it with the former,
though user may choose th
### Requirements and Supported versions
- JDK v17
-- Debezium Connectors v3.1
-- Apache Beam 2.66
+- Debezium Connectors v3.5
+- Apache Beam 2.76
## Running Unit Tests
diff --git
a/sdks/java/io/debezium/src/main/java/org/apache/beam/io/debezium/DebeziumIO.java
b/sdks/java/io/debezium/src/main/java/org/apache/beam/io/debezium/DebeziumIO.java
index 6c31d5a0234..b89b3644f61 100644
---
a/sdks/java/io/debezium/src/main/java/org/apache/beam/io/debezium/DebeziumIO.java
+++
b/sdks/java/io/debezium/src/main/java/org/apache/beam/io/debezium/DebeziumIO.java
@@ -36,7 +36,6 @@ import org.apache.beam.sdk.values.PBegin;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Joiner;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
-import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Maps;
import org.apache.kafka.connect.source.SourceConnector;
import org.apache.kafka.connect.source.SourceRecord;
import org.checkerframework.checker.nullness.qual.Nullable;
@@ -75,7 +74,7 @@ import org.slf4j.LoggerFactory;
* .withConnectorClass(MySqlConnector.class)
* .withConnectionProperty("database.server.id", "184054")
* .withConnectionProperty("database.server.name", "serverid")
- * .withConnectionProperty("database.history",
DebeziumSDFDatabaseHistory.class.getName())
+ * .withConnectionProperty("schema.history.internal",
DebeziumSDFDatabaseHistory.class.getName())
* .withConnectionProperty("include.schema.changes", "false");
*
* PipelineOptions options = PipelineOptionsFactory.create();
@@ -313,9 +312,9 @@ public class DebeziumIO {
new KafkaSourceConsumerFn.OffsetTracker(
new KafkaSourceConsumerFn.OffsetHolder(null, null, 0)));
- Map<String, String> connectorConfig =
- Maps.newHashMap(getConnectorConfiguration().getConfigurationMap());
- connectorConfig.put("snapshot.mode", "schema_only");
+ // Deliberately runs with the connector's configured snapshot mode:
schema inference samples
+ // an actual data record, which a schema-only snapshot ("no_data",
formerly "schema_only")
+ // would never emit.
SourceRecord sampledRecord =
fn.getOneRecord(getConnectorConfiguration().getConfigurationMap());
fn.reset();
@@ -641,10 +640,10 @@ public class DebeziumIO {
configuration.computeIfAbsent(entry.getKey(), k -> entry.getValue());
}
- // Set default Database History impl. if not provided implementation and
Kafka topic prefix,
- // if not provided
+ // Set default schema history impl. if not provided implementation and
Kafka topic prefix,
+ // if not provided. Before Debezium 2.0 this key was named
"database.history".
configuration.computeIfAbsent(
- "database.history",
+ "schema.history.internal",
k ->
KafkaSourceConsumerFn.DebeziumSDFDatabaseHistory.class.getName());
configuration.computeIfAbsent("topic.prefix", k ->
"beam-debezium-connector");
configuration.computeIfAbsent(
diff --git
a/sdks/java/io/debezium/src/main/java/org/apache/beam/io/debezium/KafkaSourceConsumerFn.java
b/sdks/java/io/debezium/src/main/java/org/apache/beam/io/debezium/KafkaSourceConsumerFn.java
index d298ddd9caf..89fc2a5c085 100644
---
a/sdks/java/io/debezium/src/main/java/org/apache/beam/io/debezium/KafkaSourceConsumerFn.java
+++
b/sdks/java/io/debezium/src/main/java/org/apache/beam/io/debezium/KafkaSourceConsumerFn.java
@@ -350,6 +350,11 @@ public class KafkaSourceConsumerFn<T> extends
DoFn<Map<String, String>, T> {
LOG.debug("------------- Creating an offset storage reader");
return new DebeziumSourceOffsetStorageReader(initialOffset);
}
+
+ @Override
+ public org.apache.kafka.common.metrics.PluginMetrics pluginMetrics() {
+ return null;
+ }
}
private static class DebeziumSourceOffsetStorageReader implements
OffsetStorageReader {
diff --git
a/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumIOMySqlConnectorIT.java
b/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumIOMySqlConnectorIT.java
index 3fe86a29cce..e0f811a8495 100644
---
a/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumIOMySqlConnectorIT.java
+++
b/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumIOMySqlConnectorIT.java
@@ -74,7 +74,7 @@ public class DebeziumIOMySqlConnectorIT {
@ClassRule
public static final MySQLContainer<?> MY_SQL_CONTAINER =
new MySQLContainer<>(
-
DockerImageName.parse("quay.io/debezium/example-mysql:3.1.3.Final")
+
DockerImageName.parse("quay.io/debezium/example-mysql:3.5.2.Final")
.asCompatibleSubstituteFor("mysql"))
.withPassword("debezium")
.withUsername("mysqluser")
@@ -277,8 +277,8 @@ public class DebeziumIOMySqlConnectorIT {
.withMaxNumberOfRecords(30)
.withCoder(StringUtf8Coder.of()));
String expected =
-
"{\"metadata\":{\"connector\":\"mysql\",\"version\":\"3.1.3.Final\",\"name\":\"beam-debezium-connector\","
- +
"\"database\":\"inventory\",\"schema\":\"binlog.000002\",\"table\":\"addresses\"},\"before\":null,"
+
"{\"metadata\":{\"connector\":\"mysql\",\"version\":\"3.5.2.Final\",\"name\":\"beam-debezium-connector\","
+ +
"\"database\":\"inventory\",\"schema\":\"mysql-bin.000003\",\"table\":\"addresses\"},\"before\":null,"
+ "\"after\":{\"fields\":{\"zip\":\"76036\",\"city\":\"Euless\","
+ "\"street\":\"3183 Moore
Avenue\",\"id\":10,\"state\":\"Texas\",\"customer_id\":1001,"
+ "\"type\":\"SHIPPING\"}}}";
diff --git
a/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumIOPostgresSqlConnectorIT.java
b/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumIOPostgresSqlConnectorIT.java
index 87b9bbb92e5..0a76415d23b 100644
---
a/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumIOPostgresSqlConnectorIT.java
+++
b/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumIOPostgresSqlConnectorIT.java
@@ -56,7 +56,7 @@ public class DebeziumIOPostgresSqlConnectorIT {
@ClassRule
public static final PostgreSQLContainer<?> POSTGRES_SQL_CONTAINER =
new PostgreSQLContainer<>(
-
DockerImageName.parse("quay.io/debezium/example-postgres:3.1.3.Final")
+
DockerImageName.parse("quay.io/debezium/example-postgres:3.5.2.Final")
.asCompatibleSubstituteFor("postgres"))
.withPassword("dbz")
.withUsername("debezium")
@@ -180,7 +180,7 @@ public class DebeziumIOPostgresSqlConnectorIT {
.withMaxNumberOfRecords(30)
.withCoder(StringUtf8Coder.of()));
String expected =
-
"{\"metadata\":{\"connector\":\"postgresql\",\"version\":\"3.1.3.Final\",\"name\":\"beam-debezium-connector\","
+
"{\"metadata\":{\"connector\":\"postgresql\",\"version\":\"3.5.2.Final\",\"name\":\"beam-debezium-connector\","
+
"\"database\":\"inventory\",\"schema\":\"inventory\",\"table\":\"customers\"},\"before\":null,"
+
"\"after\":{\"fields\":{\"last_name\":\"Thomas\",\"id\":1001,\"first_name\":\"Sally\","
+ "\"email\":\"[email protected]\"}}}";
diff --git
a/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumIOTest.java
b/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumIOTest.java
index 80509f5bb91..074f0ba5b52 100644
---
a/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumIOTest.java
+++
b/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumIOTest.java
@@ -52,7 +52,8 @@ public class DebeziumIOTest implements Serializable {
.withConnectionProperty("database.server.id", "184054")
.withConnectionProperty("database.server.name", "dbserver1")
.withConnectionProperty(
- "database.history",
KafkaSourceConsumerFn.DebeziumSDFDatabaseHistory.class.getName())
+ "schema.history.internal",
+ KafkaSourceConsumerFn.DebeziumSDFDatabaseHistory.class.getName())
.withConnectionProperty("include.schema.changes", "false");
@Test
diff --git
a/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumReadSchemaTransformTest.java
b/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumReadSchemaTransformTest.java
index 2fc8996ba55..b961ad84a7b 100644
---
a/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumReadSchemaTransformTest.java
+++
b/sdks/java/io/debezium/src/test/java/org/apache/beam/io/debezium/DebeziumReadSchemaTransformTest.java
@@ -18,6 +18,7 @@
package org.apache.beam.io.debezium;
import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertThrows;
import io.debezium.DebeziumException;
@@ -60,7 +61,7 @@ public class DebeziumReadSchemaTransformTest {
@ClassRule
public static final MySQLContainer<?> MY_SQL_CONTAINER =
new MySQLContainer<>(
- DockerImageName.parse("debezium/example-mysql:1.4")
+
DockerImageName.parse("quay.io/debezium/example-mysql:3.5.2.Final")
.asCompatibleSubstituteFor("mysql"))
.withPassword("debezium")
.withUsername("mysqluser")
@@ -118,6 +119,17 @@ public class DebeziumReadSchemaTransformTest {
.build());
}
+ // Since Debezium 3.5 connection failures surface as a RetriableException
wrapping the
+ // DebeziumException instead of a top-level DebeziumException.
+ private static DebeziumException findDebeziumException(Throwable thrown) {
+ Throwable cause = thrown;
+ while (cause != null && !(cause instanceof DebeziumException)) {
+ cause = cause.getCause();
+ }
+ assertNotNull("Expected DebeziumException in cause chain", cause);
+ return (DebeziumException) cause;
+ }
+
@Test
public void testNoProblem() {
Pipeline readPipeline = Pipeline.create();
@@ -142,9 +154,9 @@ public class DebeziumReadSchemaTransformTest {
@Test
public void testWrongUser() {
Pipeline readPipeline = Pipeline.create();
- DebeziumException ex =
+ Exception thrown =
assertThrows(
- DebeziumException.class,
+ Exception.class,
() -> {
PCollectionRowTuple.empty(readPipeline)
.apply(
@@ -156,6 +168,7 @@ public class DebeziumReadSchemaTransformTest {
"localhost"))
.get("output");
});
+ DebeziumException ex = findDebeziumException(thrown);
assertThat(ex.getCause().getMessage(),
Matchers.containsString("password"));
assertThat(ex.getCause().getMessage(),
Matchers.containsString("wrongUser"));
}
@@ -163,9 +176,9 @@ public class DebeziumReadSchemaTransformTest {
@Test
public void testWrongPassword() {
Pipeline readPipeline = Pipeline.create();
- DebeziumException ex =
+ Exception thrown =
assertThrows(
- DebeziumException.class,
+ Exception.class,
() -> {
PCollectionRowTuple.empty(readPipeline)
.apply(
@@ -177,6 +190,7 @@ public class DebeziumReadSchemaTransformTest {
"localhost"))
.get("output");
});
+ DebeziumException ex = findDebeziumException(thrown);
assertThat(ex.getCause().getMessage(),
Matchers.containsString("password"));
assertThat(ex.getCause().getMessage(), Matchers.containsString(userName));
}
@@ -184,14 +198,15 @@ public class DebeziumReadSchemaTransformTest {
@Test
public void testWrongPort() {
Pipeline readPipeline = Pipeline.create();
- DebeziumException ex =
+ Exception thrown =
assertThrows(
- DebeziumException.class,
+ Exception.class,
() -> {
PCollectionRowTuple.empty(readPipeline)
.apply(makePtransform(userName, password, database, 12345,
"localhost"))
.get("output");
});
+ DebeziumException ex = findDebeziumException(thrown);
Throwable lowestCause = ex.getCause();
while (lowestCause.getCause() != null) {
lowestCause = lowestCause.getCause();
diff --git a/sdks/python/apache_beam/io/external/xlang_debeziumio_it_test.py
b/sdks/python/apache_beam/io/external/xlang_debeziumio_it_test.py
index 30b96f01a1a..5fc7c1567cd 100644
--- a/sdks/python/apache_beam/io/external/xlang_debeziumio_it_test.py
+++ b/sdks/python/apache_beam/io/external/xlang_debeziumio_it_test.py
@@ -89,7 +89,7 @@ class CrossLanguageDebeziumIOTest(unittest.TestCase):
expected_response = [{
"metadata": {
"connector": "postgresql",
- "version": "3.1.3.Final",
+ "version": "3.5.2.Final",
"name": "beam-debezium-connector",
"database": "inventory",
"schema": "inventory",