Skip to content

Commit 7a09623

Browse files
committed
Add mysql read to managed io
1 parent 1b05ebb commit 7a09623

5 files changed

Lines changed: 37 additions & 1 deletion

File tree

model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/external_transforms.proto

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -80,6 +80,8 @@ message ManagedTransforms {
8080
"beam:schematransform:org.apache.beam:postgres_read:v1"];
8181
POSTGRES_WRITE = 8 [(org.apache.beam.model.pipeline.v1.beam_urn) =
8282
"beam:schematransform:org.apache.beam:postgres_write:v1"];
83+
MYSQL_READ = 9 [(org.apache.beam.model.pipeline.v1.beam_urn) =
84+
"beam:schematransform:org.apache.beam:mysql_read:v1"];
8385
}
8486
}
8587

sdks/java/io/jdbc/src/main/java/org/apache/beam/sdk/io/jdbc/providers/ReadFromMySqlSchemaTransformProvider.java

Lines changed: 30 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,20 +18,28 @@
1818
package org.apache.beam.sdk.io.jdbc.providers;
1919

2020
import static org.apache.beam.sdk.io.jdbc.JdbcUtil.MYSQL;
21+
import static org.apache.beam.sdk.util.construction.BeamUrns.getUrn;
2122

2223
import com.google.auto.service.AutoService;
24+
import org.apache.beam.model.pipeline.v1.ExternalTransforms;
2325
import org.apache.beam.sdk.io.jdbc.JdbcReadSchemaTransformProvider;
26+
import org.apache.beam.sdk.schemas.transforms.SchemaTransform;
2427
import org.apache.beam.sdk.schemas.transforms.SchemaTransformProvider;
2528
import org.checkerframework.checker.initialization.qual.Initialized;
2629
import org.checkerframework.checker.nullness.qual.NonNull;
2730
import org.checkerframework.checker.nullness.qual.UnknownKeyFor;
31+
import org.slf4j.Logger;
32+
import org.slf4j.LoggerFactory;
2833

2934
@AutoService(SchemaTransformProvider.class)
3035
public class ReadFromMySqlSchemaTransformProvider extends JdbcReadSchemaTransformProvider {
3136

37+
private static final Logger LOG =
38+
LoggerFactory.getLogger(ReadFromMySqlSchemaTransformProvider.class);
39+
3240
@Override
3341
public @UnknownKeyFor @NonNull @Initialized String identifier() {
34-
return "beam:schematransform:org.apache.beam:mysql_read:v1";
42+
return getUrn(ExternalTransforms.ManagedTransforms.Urns.MYSQL_READ);
3543
}
3644

3745
@Override
@@ -43,4 +51,25 @@ public String description() {
4351
protected String jdbcType() {
4452
return MYSQL;
4553
}
54+
55+
@Override
56+
public @UnknownKeyFor @NonNull @Initialized SchemaTransform from(
57+
JdbcReadSchemaTransformConfiguration configuration) {
58+
String jdbcType = configuration.getJdbcType();
59+
if (jdbcType != null && !jdbcType.equals(jdbcType())) {
60+
throw new IllegalArgumentException(
61+
String.format("Wrong JDBC type. Expected '%s' but got '%s'", jdbcType(), jdbcType));
62+
}
63+
64+
Integer fetchSize = configuration.getFetchSize();
65+
if (fetchSize != null
66+
&& fetchSize > 0
67+
&& configuration.getJdbcUrl() != null
68+
&& !configuration.getJdbcUrl().contains("useCursorFetch=true")) {
69+
LOG.warn(
70+
"The fetchSize option is ignored. It is required to set useCursorFetch=true"
71+
+ " in the JDBC URL when using fetchSize for MySQL");
72+
}
73+
return super.from(configuration);
74+
}
4675
}

sdks/java/managed/src/main/java/org/apache/beam/sdk/managed/Managed.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -97,6 +97,7 @@ public class Managed {
9797
public static final String KAFKA = "kafka";
9898
public static final String BIGQUERY = "bigquery";
9999
public static final String POSTGRES = "postgres";
100+
public static final String MYSQL = "mysql";
100101

101102
// Supported SchemaTransforms
102103
public static final Map<String, String> READ_TRANSFORMS =
@@ -106,6 +107,7 @@ public class Managed {
106107
.put(KAFKA, getUrn(ExternalTransforms.ManagedTransforms.Urns.KAFKA_READ))
107108
.put(BIGQUERY, getUrn(ExternalTransforms.ManagedTransforms.Urns.BIGQUERY_READ))
108109
.put(POSTGRES, getUrn(ExternalTransforms.ManagedTransforms.Urns.POSTGRES_READ))
110+
.put(MYSQL, getUrn(ExternalTransforms.ManagedTransforms.Urns.MYSQL_READ))
109111
.build();
110112
public static final Map<String, String> WRITE_TRANSFORMS =
111113
ImmutableMap.<String, String>builder()

sdks/python/apache_beam/transforms/external.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -83,6 +83,7 @@
8383
ManagedTransforms.Urns.BIGQUERY_WRITE.urn: _GCP_EXPANSION_SERVICE_JAR_TARGET, # pylint: disable=line-too-long
8484
ManagedTransforms.Urns.POSTGRES_READ.urn: _GCP_EXPANSION_SERVICE_JAR_TARGET,
8585
ManagedTransforms.Urns.POSTGRES_WRITE.urn: _GCP_EXPANSION_SERVICE_JAR_TARGET, # pylint: disable=line-too-long
86+
ManagedTransforms.Urns.MYSQL_READ.urn: _GCP_EXPANSION_SERVICE_JAR_TARGET,
8687
}
8788

8889

sdks/python/apache_beam/transforms/managed.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -86,6 +86,7 @@
8686
KAFKA = "kafka"
8787
BIGQUERY = "bigquery"
8888
POSTGRES = "postgres"
89+
MYSQL = "mysql"
8990

9091
__all__ = ["ICEBERG", "KAFKA", "BIGQUERY", "Read", "Write"]
9192

@@ -98,6 +99,7 @@ class Read(PTransform):
9899
KAFKA: ManagedTransforms.Urns.KAFKA_READ.urn,
99100
BIGQUERY: ManagedTransforms.Urns.BIGQUERY_READ.urn,
100101
POSTGRES: ManagedTransforms.Urns.POSTGRES_READ.urn,
102+
MYSQL: ManagedTransforms.Urns.MYSQL_READ.urn,
101103
}
102104

103105
def __init__(

0 commit comments

Comments
 (0)