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.

Reply via email to