This is an automated email from the ASF dual-hosted git repository.
tvalentyn 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 4d45c70ecc7 Several small updates to YAML ML examples (#40185)
4d45c70ecc7 is described below
commit 4d45c70ecc78ebad7cc1c7f1e0de7c13f3181666
Author: Chamikara Jayalath <[email protected]>
AuthorDate: Mon Sep 21 11:49:43 2026 -0700
Several small updates to YAML ML examples (#40185)
---
.../yaml/examples/transforms/ml/enrich_spanner_with_bigquery.yaml | 3 ++-
.../yaml/examples/transforms/ml/log_analysis/ml_preprocessing.yaml | 4 +++-
.../apache_beam/yaml/examples/transforms/ml/log_analysis/train.py | 4 +++-
3 files changed, 8 insertions(+), 3 deletions(-)
diff --git
a/sdks/python/apache_beam/yaml/examples/transforms/ml/enrich_spanner_with_bigquery.yaml
b/sdks/python/apache_beam/yaml/examples/transforms/ml/enrich_spanner_with_bigquery.yaml
index e63b3105cc0..e9c140a7965 100644
---
a/sdks/python/apache_beam/yaml/examples/transforms/ml/enrich_spanner_with_bigquery.yaml
+++
b/sdks/python/apache_beam/yaml/examples/transforms/ml/enrich_spanner_with_bigquery.yaml
@@ -36,7 +36,8 @@ pipeline:
handler_config:
project: "apache-beam-testing"
table_name: "apache-beam-testing.ALL_TEST.customers"
- row_restriction_template: "customer_id = 1001 or customer_id = 1003"
+ # Use placeholder format string so values from 'fields' dynamically
populate the WHERE clause.
+ row_restriction_template: "customer_id = {}"
fields: ["customer_id"]
# Step 3: Map enriched values to Beam schema
diff --git
a/sdks/python/apache_beam/yaml/examples/transforms/ml/log_analysis/ml_preprocessing.yaml
b/sdks/python/apache_beam/yaml/examples/transforms/ml/log_analysis/ml_preprocessing.yaml
index e567a46476b..08b2dacf496 100644
---
a/sdks/python/apache_beam/yaml/examples/transforms/ml/log_analysis/ml_preprocessing.yaml
+++
b/sdks/python/apache_beam/yaml/examples/transforms/ml/log_analysis/ml_preprocessing.yaml
@@ -117,7 +117,9 @@ pipeline:
options:
yaml_experimental_features: [ 'ML' ]
- temp_location: "gs://apache-beam-testing/temp"
+ # Dynamic temp location avoids hardcoding inaccessible staging buckets
across environments.
+ temp_location: "{{ WAREHOUSE }}/temp"
+
# Expected:
# Row(id=1, date='2024-10-01', time='12:00:00', level='INFO', process='Main',
component='ComponentA', content='System started successfully',
embedding=[0.13483997249264842, 0.26967994498529685, 0.40451991747794525,
0.5393598899705937, 0.674199862463242])
diff --git
a/sdks/python/apache_beam/yaml/examples/transforms/ml/log_analysis/train.py
b/sdks/python/apache_beam/yaml/examples/transforms/ml/log_analysis/train.py
index f0f957aa7ba..ce982008de8 100644
--- a/sdks/python/apache_beam/yaml/examples/transforms/ml/log_analysis/train.py
+++ b/sdks/python/apache_beam/yaml/examples/transforms/ml/log_analysis/train.py
@@ -55,7 +55,9 @@ class ModelHelper():
def load_data(self):
logging.info("Querying vector embeddings from BigQuery...")
- client = bigquery.Client()
+ # Extract project from table spec to avoid failures when no default
project configured.
+ project = self.bq_table.split('.')[0] if '.' in self.bq_table else None
+ client = bigquery.Client(project=project)
sql = f"""
SELECT *
FROM `{self.bq_table}`