Skip to content

Commit f86452f

Browse files
authored
Support managed jdbc io (MySQL) (#36045)
* Add mysql read to managed io * Add mysql write to managed io * Add schema transform translation and test for mysql read and write * Remove redundant config validation. * Allow jdbcType to be empty. * Address reviewer's comments.
1 parent 9a5a357 commit f86452f

8 files changed

Lines changed: 406 additions & 2 deletions

File tree

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

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -80,6 +80,10 @@ 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"];
85+
MYSQL_WRITE = 10 [(org.apache.beam.model.pipeline.v1.beam_urn) =
86+
"beam:schematransform:org.apache.beam:mysql_write:v1"];
8387
}
8488
}
8589

Lines changed: 93 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,93 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing, software
13+
* distributed under the License is distributed on an "AS IS" BASIS,
14+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15+
* See the License for the specific language governing permissions and
16+
* limitations under the License.
17+
*/
18+
package org.apache.beam.sdk.io.jdbc.providers;
19+
20+
import static org.apache.beam.sdk.io.jdbc.providers.ReadFromMySqlSchemaTransformProvider.MySqlReadSchemaTransform;
21+
import static org.apache.beam.sdk.io.jdbc.providers.WriteToMySqlSchemaTransformProvider.MySqlWriteSchemaTransform;
22+
import static org.apache.beam.sdk.schemas.transforms.SchemaTransformTranslation.SchemaTransformPayloadTranslator;
23+
24+
import com.google.auto.service.AutoService;
25+
import java.util.Map;
26+
import org.apache.beam.sdk.schemas.transforms.SchemaTransformProvider;
27+
import org.apache.beam.sdk.transforms.PTransform;
28+
import org.apache.beam.sdk.util.construction.PTransformTranslation;
29+
import org.apache.beam.sdk.util.construction.TransformPayloadTranslatorRegistrar;
30+
import org.apache.beam.sdk.values.Row;
31+
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
32+
33+
public class MySqlSchemaTransformTranslation {
34+
static class MySqlReadSchemaTransformTranslator
35+
extends SchemaTransformPayloadTranslator<MySqlReadSchemaTransform> {
36+
@Override
37+
public SchemaTransformProvider provider() {
38+
return new ReadFromMySqlSchemaTransformProvider();
39+
}
40+
41+
@Override
42+
public Row toConfigRow(MySqlReadSchemaTransform transform) {
43+
return transform.getConfigurationRow();
44+
}
45+
}
46+
47+
@AutoService(TransformPayloadTranslatorRegistrar.class)
48+
public static class ReadRegistrar implements TransformPayloadTranslatorRegistrar {
49+
@Override
50+
@SuppressWarnings({
51+
"rawtypes",
52+
})
53+
public Map<
54+
? extends Class<? extends PTransform>,
55+
? extends PTransformTranslation.TransformPayloadTranslator>
56+
getTransformPayloadTranslators() {
57+
return ImmutableMap
58+
.<Class<? extends PTransform>, PTransformTranslation.TransformPayloadTranslator>builder()
59+
.put(MySqlReadSchemaTransform.class, new MySqlReadSchemaTransformTranslator())
60+
.build();
61+
}
62+
}
63+
64+
static class MySqlWriteSchemaTransformTranslator
65+
extends SchemaTransformPayloadTranslator<MySqlWriteSchemaTransform> {
66+
@Override
67+
public SchemaTransformProvider provider() {
68+
return new WriteToMySqlSchemaTransformProvider();
69+
}
70+
71+
@Override
72+
public Row toConfigRow(MySqlWriteSchemaTransform transform) {
73+
return transform.getConfigurationRow();
74+
}
75+
}
76+
77+
@AutoService(TransformPayloadTranslatorRegistrar.class)
78+
public static class WriteRegistrar implements TransformPayloadTranslatorRegistrar {
79+
@Override
80+
@SuppressWarnings({
81+
"rawtypes",
82+
})
83+
public Map<
84+
? extends Class<? extends PTransform>,
85+
? extends PTransformTranslation.TransformPayloadTranslator>
86+
getTransformPayloadTranslators() {
87+
return ImmutableMap
88+
.<Class<? extends PTransform>, PTransformTranslation.TransformPayloadTranslator>builder()
89+
.put(MySqlWriteSchemaTransform.class, new MySqlWriteSchemaTransformTranslator())
90+
.build();
91+
}
92+
}
93+
}

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

Lines changed: 40 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,35 @@ 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.isEmpty() && !jdbcType.equals(jdbcType())) {
60+
LOG.warn(
61+
"Wrong JDBC type. Expected '{}' but got '{}'. Overriding with '{}'.",
62+
jdbcType(),
63+
jdbcType,
64+
jdbcType());
65+
configuration = configuration.toBuilder().setJdbcType(jdbcType()).build();
66+
}
67+
68+
Integer fetchSize = configuration.getFetchSize();
69+
if (fetchSize != null
70+
&& fetchSize > 0
71+
&& configuration.getJdbcUrl() != null
72+
&& !configuration.getJdbcUrl().contains("useCursorFetch=true")) {
73+
throw new IllegalArgumentException(
74+
"It is required to set useCursorFetch=true"
75+
+ " in the JDBC URL when using fetchSize for MySQL");
76+
}
77+
return new MySqlReadSchemaTransform(configuration);
78+
}
79+
80+
public static class MySqlReadSchemaTransform extends JdbcReadSchemaTransform {
81+
public MySqlReadSchemaTransform(JdbcReadSchemaTransformConfiguration config) {
82+
super(config, MYSQL);
83+
}
84+
}
4685
}

sdks/java/io/jdbc/src/main/java/org/apache/beam/sdk/io/jdbc/providers/WriteToMySqlSchemaTransformProvider.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.JdbcWriteSchemaTransformProvider;
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 WriteToMySqlSchemaTransformProvider extends JdbcWriteSchemaTransformProvider {
3136

37+
private static final Logger LOG =
38+
LoggerFactory.getLogger(WriteToMySqlSchemaTransformProvider.class);
39+
3240
@Override
3341
public @UnknownKeyFor @NonNull @Initialized String identifier() {
34-
return "beam:schematransform:org.apache.beam:mysql_write:v1";
42+
return getUrn(ExternalTransforms.ManagedTransforms.Urns.MYSQL_WRITE);
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+
JdbcWriteSchemaTransformConfiguration configuration) {
58+
String jdbcType = configuration.getJdbcType();
59+
if (jdbcType != null && !jdbcType.isEmpty() && !jdbcType.equals(jdbcType())) {
60+
LOG.warn(
61+
"Wrong JDBC type. Expected '{}' but got '{}'. Overriding with '{}'.",
62+
jdbcType(),
63+
jdbcType,
64+
jdbcType());
65+
configuration = configuration.toBuilder().setJdbcType(jdbcType()).build();
66+
}
67+
return new MySqlWriteSchemaTransform(configuration);
68+
}
69+
70+
public static class MySqlWriteSchemaTransform extends JdbcWriteSchemaTransform {
71+
public MySqlWriteSchemaTransform(JdbcWriteSchemaTransformConfiguration config) {
72+
super(config, MYSQL);
73+
}
74+
}
4675
}

0 commit comments

Comments
 (0)