This is an automated email from the ASF dual-hosted git repository.
jrmccluskey 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 274905845dd Add examples for cloudsql (#40027)
274905845dd is described below
commit 274905845ddc538fbf44af3cfcc2c9b987179c67
Author: Derrick Williams <[email protected]>
AuthorDate: Wed Oct 7 11:38:47 2026 -0400
Add examples for cloudsql (#40027)
* Add examples for cloudsql
* remove changes
* remove advanced sections
---
.../io/jdbc/JdbcReadSchemaTransformProvider.java | 8 +--
.../io/jdbc/JdbcWriteSchemaTransformProvider.java | 8 +--
.../yaml/examples/testing/examples_test.py | 2 +
..._to_bigquery.yaml => cloudsql_to_bigquery.yaml} | 22 +++---
.../transforms/blueprint/postgres_to_bigquery.yaml | 5 +-
.../extended_tests/e2e/cloudsql_to_bigquery.yaml | 82 ++++++++++++++++++++++
sdks/python/apache_beam/yaml/integration_tests.py | 35 +++++++++
7 files changed, 131 insertions(+), 31 deletions(-)
diff --git
a/sdks/java/io/jdbc/src/main/java/org/apache/beam/sdk/io/jdbc/JdbcReadSchemaTransformProvider.java
b/sdks/java/io/jdbc/src/main/java/org/apache/beam/sdk/io/jdbc/JdbcReadSchemaTransformProvider.java
index 9b29d76673d..c1417aa11c5 100644
---
a/sdks/java/io/jdbc/src/main/java/org/apache/beam/sdk/io/jdbc/JdbcReadSchemaTransformProvider.java
+++
b/sdks/java/io/jdbc/src/main/java/org/apache/beam/sdk/io/jdbc/JdbcReadSchemaTransformProvider.java
@@ -152,13 +152,7 @@ public class JdbcReadSchemaTransformProvider
+ " - type: %s%n"
+ " config:%n"
+ " url: \"jdbc:%s://my-host:%d/database\"%n"
- + " table: \"my-table\"%n"
- + "%n"
- + "#### Advanced Usage%n"
- + "%n"
- + "It might be necessary to use a custom JDBC driver that is not
packaged with this "
- + "transform. If that is the case, see ReadFromJdbc which "
- + "allows for more custom configuration.",
+ + " table: \"my-table\"%n",
prettyName,
prettyName,
transformName,
diff --git
a/sdks/java/io/jdbc/src/main/java/org/apache/beam/sdk/io/jdbc/JdbcWriteSchemaTransformProvider.java
b/sdks/java/io/jdbc/src/main/java/org/apache/beam/sdk/io/jdbc/JdbcWriteSchemaTransformProvider.java
index 3379fad639e..02bd132e704 100644
---
a/sdks/java/io/jdbc/src/main/java/org/apache/beam/sdk/io/jdbc/JdbcWriteSchemaTransformProvider.java
+++
b/sdks/java/io/jdbc/src/main/java/org/apache/beam/sdk/io/jdbc/JdbcWriteSchemaTransformProvider.java
@@ -156,13 +156,7 @@ public class JdbcWriteSchemaTransformProvider
+ " - type: %s%n"
+ " config:%n"
+ " url: \"jdbc:%s://my-host:%d/database\"%n"
- + " table: \"my-table\"%n"
- + "%n"
- + "#### Advanced Usage%n"
- + "%n"
- + "It might be necessary to use a custom JDBC driver that is not
packaged with this "
- + "transform. If that is the case, see WriteToJdbc which "
- + "allows for more custom configuration.",
+ + " table: \"my-table\"%n",
prettyName,
prettyName,
transformName,
diff --git a/sdks/python/apache_beam/yaml/examples/testing/examples_test.py
b/sdks/python/apache_beam/yaml/examples/testing/examples_test.py
index 97acf427f54..3ea7e0f07d8 100644
--- a/sdks/python/apache_beam/yaml/examples/testing/examples_test.py
+++ b/sdks/python/apache_beam/yaml/examples/testing/examples_test.py
@@ -669,6 +669,7 @@ def _kafka_test_preprocessor(
'test_gcs_text_to_bigquery_yaml',
'test_sqlserver_to_bigquery_yaml',
'test_postgres_to_bigquery_yaml',
+ 'test_cloudsql_to_bigquery_yaml',
'test_kafka_to_iceberg_yaml',
'test_pubsub_to_iceberg_yaml',
'test_oracle_to_bigquery_yaml',
@@ -931,6 +932,7 @@ def __sqlserver_io_read_test_preprocessor(
@YamlExamplesTestSuite.register_test_preprocessor([
'test_postgres_to_bigquery_yaml',
+ 'test_cloudsql_to_bigquery_yaml',
])
def __postgres_io_read_test_preprocessor(
test_spec: dict, expected: list[str], env: TestEnvironment):
diff --git
a/sdks/python/apache_beam/yaml/examples/transforms/blueprint/postgres_to_bigquery.yaml
b/sdks/python/apache_beam/yaml/examples/transforms/blueprint/cloudsql_to_bigquery.yaml
similarity index 75%
copy from
sdks/python/apache_beam/yaml/examples/transforms/blueprint/postgres_to_bigquery.yaml
copy to
sdks/python/apache_beam/yaml/examples/transforms/blueprint/cloudsql_to_bigquery.yaml
index b532636f46e..ee8885f2815 100644
---
a/sdks/python/apache_beam/yaml/examples/transforms/blueprint/postgres_to_bigquery.yaml
+++
b/sdks/python/apache_beam/yaml/examples/transforms/blueprint/cloudsql_to_bigquery.yaml
@@ -15,36 +15,30 @@
# See the License for the specific language governing permissions and
# limitations under the License.
-# This is an example of a Beam YAML pipeline that reads from spanner database
-# and writes to GCS avro files. This matches the Dataflow Template located
-# here -
https://cloud.google.com/dataflow/docs/guides/templates/provided/cloud-spanner-to-avro
+# This is an example of a Beam YAML pipeline that reads from a Google Cloud SQL
+# PostgreSQL database using the Cloud SQL JDBC Socket Factory and writes to
BigQuery.
pipeline:
type: composite
transforms:
- # Step 1: Reading data from Postgres
+ # Step 1: Reading data from Cloud SQL Postgres
- type: ReadFromPostgres
- name: ReadFromPostgres
+ name: ReadFromCloudSqlPostgres
config:
- url:
"jdbc:postgresql://localhost:12345/shipment?user=user&password=postgres123"
+ url:
"jdbc:postgresql:///shipment?cloudSqlInstance=my-project:us-central1:my-instance&socketFactory=com.google.cloud.sql.postgres.SocketFactory"
query: "SELECT * FROM shipments"
driver_class_name: "org.postgresql.Driver"
+ username: "my-username"
+ password: "my-password"
# Step 2: Write records out to BigQuery
- type: WriteToBigQuery
name: WriteShipments
- input: ReadFromPostgres
+ input: ReadFromCloudSqlPostgres
config:
table: "apache-beam-testing.yaml_test.shipments"
create_disposition: "CREATE_NEVER"
write_disposition: "WRITE_APPEND"
- error_handling:
- output: "deadLetterQueue"
num_streams: 1
- # Step 3: Write the failed messages to BQ to a dead letter queue JSON file
- - type: WriteToJson
- input: WriteShipments.deadLetterQueue
- config:
- path: "gs://my-bucket/yaml-123/writingToBigQueryErrors.json"
options:
temp_location: "gs://apache-beam-testing/temp"
diff --git
a/sdks/python/apache_beam/yaml/examples/transforms/blueprint/postgres_to_bigquery.yaml
b/sdks/python/apache_beam/yaml/examples/transforms/blueprint/postgres_to_bigquery.yaml
index b532636f46e..4ca4cd2461c 100644
---
a/sdks/python/apache_beam/yaml/examples/transforms/blueprint/postgres_to_bigquery.yaml
+++
b/sdks/python/apache_beam/yaml/examples/transforms/blueprint/postgres_to_bigquery.yaml
@@ -15,9 +15,8 @@
# See the License for the specific language governing permissions and
# limitations under the License.
-# This is an example of a Beam YAML pipeline that reads from spanner database
-# and writes to GCS avro files. This matches the Dataflow Template located
-# here -
https://cloud.google.com/dataflow/docs/guides/templates/provided/cloud-spanner-to-avro
+# This is an example of a Beam YAML pipeline that reads from a PostgreSQL
database
+# and writes to BigQuery.
pipeline:
type: composite
diff --git
a/sdks/python/apache_beam/yaml/extended_tests/e2e/cloudsql_to_bigquery.yaml
b/sdks/python/apache_beam/yaml/extended_tests/e2e/cloudsql_to_bigquery.yaml
new file mode 100644
index 00000000000..a409359626d
--- /dev/null
+++ b/sdks/python/apache_beam/yaml/extended_tests/e2e/cloudsql_to_bigquery.yaml
@@ -0,0 +1,82 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+fixtures:
+ - name: CLOUDSQL
+ type: "apache_beam.yaml.integration_tests.cloudsql_postgres_fixture"
+ - name: BQ_TABLE
+ type: "apache_beam.yaml.integration_tests.temp_bigquery_table"
+ config:
+ project: "apache-beam-testing"
+ - name: TEMP_DIR
+ type: "apache_beam.yaml.integration_tests.gcs_temp_dir"
+ config:
+ bucket: "gs://temp-storage-for-end-to-end-tests/temp-it"
+
+pipelines:
+ # Pipeline 1: Write records to Cloud SQL Postgres
+ - pipeline:
+ type: chain
+ transforms:
+ - type: Create
+ config:
+ elements:
+ - {id: 1, name: "Alice", score: 95.5}
+ - {id: 2, name: "Bob", score: 88.0}
+ - {id: 3, name: "Charlie", score: 92.3}
+ - type: WriteToPostgres
+ config:
+ url: "{CLOUDSQL[JDBC_URL]}"
+ driver_class_name: "{CLOUDSQL[DRIVER_CLASS_NAME]}"
+ username: "{CLOUDSQL[USERNAME]}"
+ password: "{CLOUDSQL[PASSWORD]}"
+ table: "{CLOUDSQL[TABLE]}"
+
+ # Pipeline 2: Read from Cloud SQL Postgres and write to BigQuery
+ - pipeline:
+ type: chain
+ transforms:
+ - type: ReadFromPostgres
+ config:
+ url: "{CLOUDSQL[JDBC_URL]}"
+ driver_class_name: "{CLOUDSQL[DRIVER_CLASS_NAME]}"
+ username: "{CLOUDSQL[USERNAME]}"
+ password: "{CLOUDSQL[PASSWORD]}"
+ query: "SELECT id, name, score FROM {CLOUDSQL[TABLE]}"
+ - type: WriteToBigQuery
+ config:
+ table: "{BQ_TABLE}"
+ options:
+ project: "apache-beam-testing"
+ temp_location: "{TEMP_DIR}"
+
+ # Pipeline 3: Read from BigQuery and verify records
+ - pipeline:
+ type: chain
+ transforms:
+ - type: ReadFromBigQuery
+ config:
+ table: "{BQ_TABLE}"
+ - type: AssertEqual
+ config:
+ elements:
+ - {id: 1, name: "Alice", score: 95.5}
+ - {id: 2, name: "Bob", score: 88.0}
+ - {id: 3, name: "Charlie", score: 92.3}
+ options:
+ project: "apache-beam-testing"
+ temp_location: "{TEMP_DIR}"
diff --git a/sdks/python/apache_beam/yaml/integration_tests.py
b/sdks/python/apache_beam/yaml/integration_tests.py
index d1a03f39fc7..54be607eaba 100644
--- a/sdks/python/apache_beam/yaml/integration_tests.py
+++ b/sdks/python/apache_beam/yaml/integration_tests.py
@@ -794,6 +794,41 @@ def temp_postgres_database_with_secret_manager(
_LOGGER.warning("Could not delete GCP secret %s: %s", secret_path, err)
[email protected]
+def cloudsql_postgres_fixture():
+ """Context manager to provide a PostgreSQL testcontainer database for
testing."""
+ default_port = 5432
+ with PostgresContainer(port=default_port) as postgres_container:
+ try:
+ engine =
sqlalchemy.create_engine(postgres_container.get_connection_url())
+ with engine.begin() as connection:
+ connection.execute(
+ sqlalchemy.text(
+ "CREATE TABLE tmp_table (id INTEGER, name VARCHAR(255), score
FLOAT);"
+ ))
+
+ jdbc_url = (
+ f"jdbc:postgresql://{postgres_container.get_container_host_ip()}:"
+ f"{postgres_container.get_exposed_port(default_port)}/"
+ f"{postgres_container.dbname}?"
+ f"user={postgres_container.username}&"
+ f"password={postgres_container.password}")
+
+ yield {
+ 'JDBC_URL': jdbc_url,
+ 'DRIVER_CLASS_NAME': 'org.postgresql.Driver',
+ 'USERNAME': postgres_container.username,
+ 'PASSWORD': postgres_container.password,
+ 'DATABASE': postgres_container.dbname,
+ 'TABLE': 'tmp_table',
+ }
+ except (psycopg2.Error, Exception) as err:
+ logging.error(
+ "Error interacting with temporary Postgres DB in
cloudsql_postgres_fixture: %s",
+ err)
+ raise err
+
+
@contextlib.contextmanager
def temp_sqlserver_database():
"""Context manager to provide a temporary SQL Server database for testing.