如何用Python实现MySQL到Google BigQuery的最优数据同步(覆盖更新)
Hey there! Let’s walk through the best Python-based solutions for incremental sync from MySQL to BigQuery, since you already have a full 3M-row dataset and need to keep it updated with daily changes (inserts, updates, deletes). I’ll break down options from easiest to most robust, so you can pick what fits your use case.
This is the go-to approach if your MySQL tables have create_time or update_time columns (which they almost always should). The idea is to only sync rows that have changed since your last successful sync.
How to implement:
- Track sync state: Create a small control table in BigQuery (e.g.,
sync_control) to store the last successful sync timestamp. This ensures you don’t reprocess data if your script fails mid-run. - Fetch incremental data: Query MySQL for rows where
update_time > last_sync_timestamp. - Merge into BigQuery: Write the incremental data to a temporary BigQuery table, then use a
MERGEstatement to upsert (update existing rows, insert new ones) into your main table. - Update sync state: After a successful sync, update the control table with the latest
update_timefrom your incremental data.
Example code snippet:
import pandas as pd from google.cloud import bigquery import mysql.connector # Initialize clients bq_client = bigquery.Client(project="your-gcp-project") mysql_conn = mysql.connector.connect( host="your-mysql-host", user="mysql-user", password="mysql-pass", database="your-db" ) # Step 1: Get last sync timestamp from control table control_query = "SELECT last_sync_ts FROM your-project.your-dataset.sync_control LIMIT 1" last_sync_ts = bq_client.query(control_query).to_dataframe().iloc[0]["last_sync_ts"] # Step 2: Pull incremental data from MySQL mysql_query = f""" SELECT * FROM your-mysql-table WHERE update_time > '{last_sync_ts}' ORDER BY update_time ASC """ incremental_df = pd.read_sql(mysql_query, mysql_conn) if not incremental_df.empty: # Step 3: Write to temp BigQuery table temp_table_id = "your-project.your-dataset.temp_incremental" incremental_df.to_gbq( temp_table_id, project_id="your-gcp-project", if_exists="replace" ) # Step 4: Merge temp data into main table merge_query = f""" MERGE INTO your-project.your-dataset.main_table t USING {temp_table_id} s ON t.id = s.id # Use your primary key here WHEN MATCHED THEN UPDATE SET col1 = s.col1, col2 = s.col2, update_time = s.update_time WHEN NOT MATCHED THEN INSERT (id, col1, col2, create_time, update_time) VALUES (s.id, s.col1, s.col2, s.create_time, s.update_time) """ bq_client.query(merge_query).result() # Step 5: Update sync timestamp new_last_sync_ts = incremental_df["update_time"].max() update_control_query = f""" UPDATE your-project.your-dataset.sync_control SET last_sync_ts = '{new_last_sync_ts}' """ bq_client.query(update_control_query).result() # Cleanup mysql_conn.close()
Pros & Cons:
- ✅ Simple to implement, no MySQL config changes needed
- ✅ Works for most business tables with timestamp fields
- ❌ Can’t handle hard deletes (unless you use soft deletes like
is_deleted=1) - ❌ Risk of data loss if
update_timeisn’t properly updated on changes
If you need to capture all changes—including hard deletes—and require strict data consistency, CDC via MySQL’s binlog is the way to go. This approach listens directly to MySQL’s transaction log to capture every insert, update, and delete in real-time.
Prerequisites:
- Enable MySQL binlog: Update your
my.cnf/my.iniwith:log-bin=mysql-bin server-id=1 # Unique ID for this MySQL instance binlog_format=ROW # Required to capture row-level changes - Restart MySQL to apply changes.
How to implement with Python:
Use the pymysqlreplication library to stream binlog events, then apply those changes to BigQuery.
Example code snippet:
from pymysqlreplication import BinLogStreamReader from pymysqlreplication.row_event import ( DeleteRowsEvent, UpdateRowsEvent, WriteRowsEvent ) from google.cloud import bigquery bq_client = bigquery.Client(project="your-gcp-project") main_table_id = "your-project.your-dataset.main_table" # Initialize binlog stream reader stream = BinLogStreamReader( connection_settings={ "host": "your-mysql-host", "port": 3306, "user": "mysql-user", "passwd": "mysql-pass" }, server_id=100, # Unique ID (different from MySQL's server-id) blocking=True, # Wait for new events only_events=[DeleteRowsEvent, UpdateRowsEvent, WriteRowsEvent] ) for binlog_event in stream: for row in binlog_event.rows: # Handle each event type if isinstance(binlog_event, WriteRowsEvent): # Insert new row into BigQuery bq_client.insert_rows_json(main_table_id, [row["values"]]) elif isinstance(binlog_event, UpdateRowsEvent): # Update existing row update_query = f""" UPDATE {main_table_id} SET col1 = '{row["after_values"]["col1"]}', col2 = '{row["after_values"]["col2"]}' WHERE id = {row["before_values"]["id"]} """ bq_client.query(update_query).result() elif isinstance(binlog_event, DeleteRowsEvent): # Delete row from BigQuery delete_query = f""" DELETE FROM {main_table_id} WHERE id = {row["values"]["id"]} """ bq_client.query(delete_query).result() # Don't forget to save the binlog position on shutdown to avoid reprocessing stream.close()
Pros & Cons:
- ✅ Captures every change (including hard deletes)
- ✅ Real-time sync capability
- ✅ Highest data consistency
- ❌ Requires MySQL config changes and ongoing maintenance
- ❌ More complex code, especially for handling edge cases (e.g., schema changes)
If you want to minimize custom code, Google’s BigQuery Data Transfer Service can handle MySQL syncs natively—including incremental updates. You can use Python to configure and trigger transfers via the Google Cloud API.
How to implement:
Use the google-cloud-bigquery-datatransfer library to create a MySQL transfer configuration, set it to sync incrementally (based on a timestamp column), and schedule daily runs.
Key Benefits:
- ✅ Managed service, no need to maintain sync scripts
- ✅ Built-in error handling and retries
- ✅ Supports full initial sync + incremental updates
- ❌ Less customization for complex business logic
- Batch processing: Always handle data in batches (e.g., 10k rows at a time) to avoid rate limits and improve performance.
- Retry logic: Use libraries like
tenacityto retry failed API calls (BigQuery/MySQL) due to network issues or throttling. - Data validation: After each sync, compare row counts between MySQL’s incremental data and BigQuery to catch discrepancies.
- Logging: Implement detailed logging to track sync times, row counts, errors, and binlog positions (for CDC).
- Idempotency: Ensure your sync logic can be rerun without causing duplicates or data corruption (this is where
MERGEand binlog position tracking shine).
内容的提问来源于stack exchange,提问作者Matu Mlkz

