Airflow DAG从MySQL同步到Cloud Storage/BigQuery遇UnicodeDecodeError求助
Let’s tackle this decoding error head-on—your 0x96 byte issue is a classic case of mismatched character encodings. That byte corresponds to an en dash in Windows-1252 encoding, which isn’t valid UTF-8. Here’s how to resolve this, including configuring your Cloud SQL connection properly for Airflow:
Step 1: Configure Cloud SQL Connection String in Airflow
First, ensure your Airflow MySQL connection explicitly uses utf8mb4 (which supports all Unicode characters, including emojis and special symbols) in the connection URL.
In Airflow's UI:
- Go to Admin > Connections
- Find your
mysql_connectionentry - Update the
SQL Alchemy Connfield to include the charset parameter:mysql://[username]:[password]@[cloud-sql-ip-or-hostname]:3306/podiodb?charset=utf8mb4 - Save the connection.
This tells the MySQL driver to fetch data using utf8mb4 encoding, which should handle most special characters. However, since you’re seeing 0x96 (a Windows-1252 character), we’ll need an extra step to clean up legacy encoding.
Step 2: Fix Encoding in the Data Export Process
The default MySqlToGoogleCloudStorageOperator in Airflow 1.8 assumes all string data is UTF-8 encoded. To handle Windows-1252 characters like 0x96, we’ll create a custom operator that explicitly decodes bytes with the correct encoding.
Add this custom operator to your DAG file:
from airflow.contrib.operators.mysql_to_gcs import MySqlToGoogleCloudStorageOperator import json class MySqlToGCSFixedOperator(MySqlToGoogleCloudStorageOperator): def _write_local_data_files(self, cursor): schema = [desc[0] for desc in cursor.description] row_counter = 0 files_to_upload = [] tmp_file_handle = None for row in cursor: if row_counter % self.export_batch_size == 0: if tmp_file_handle: tmp_file_handle.close() files_to_upload.append(tmp_file_handle.name) filename = self.filename.format(row_counter) tmp_file_handle = open(filename, 'w') row_dict = dict(zip(schema, row)) # Clean up string fields with Windows-1252 encoding for key, value in row_dict.items(): if isinstance(value, str): # Decode from Windows-1252 to Unicode, replace invalid chars row_dict[key] = value.decode('windows-1252', 'replace') json.dump(row_dict, tmp_file_handle) tmp_file_handle.write('\n') row_counter += 1 if tmp_file_handle: tmp_file_handle.close() files_to_upload.append(tmp_file_handle.name) return files_to_upload
Then replace the original MySqlToGoogleCloudStorageOperator with this custom one in your DAG loop:
extract = MySqlToGCSFixedOperator( task_id="extract_mysql_%s_%s"%(connection,table), mysql_conn_id=connection, google_cloud_storage_conn_id='gcp_connection', sql="SELECT *, '%s' as source FROM podiodb.%s"%(connection,table), bucket='podio-reader-storage', filename="%s/%s/%s{}.json"%(connection,table,table), schema_filename="%s/schemas/%s.json"%(connection,table), dag=dag)
Step 3: Verify the Fix
- Test with a small dataset: Run a test query on the problematic table to isolate rows with
0x96using:
Check if the custom operator correctly converts this byte to the proper en dash character.SELECT * FROM podiodb.finance_banking_details WHERE HEX(your_column) LIKE '%96%'; - Run a single task: Trigger just the extract task for the problematic table to confirm it no longer throws the
UnicodeDecodeError. - Full DAG run: Once the extract works, run the full DAG to ensure data loads correctly into BigQuery.
Why This Works
- The connection string ensures MySQL sends data in
utf8mb4, which covers all Unicode characters. - The custom operator handles legacy Windows-1252 characters by explicitly decoding bytes and replacing any invalid sequences, preventing the JSON dump from failing.
内容的提问来源于stack exchange,提问作者Tia

