You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Airflow DAG从MySQL同步到Cloud Storage/BigQuery遇UnicodeDecodeError求助

Fixing UTF-8 Decoding Errors in Airflow MySQL-to-GCS-to-BigQuery DAG

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_connection entry
  • Update the SQL Alchemy Conn field 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

  1. Test with a small dataset: Run a test query on the problematic table to isolate rows with 0x96 using:
    SELECT * FROM podiodb.finance_banking_details WHERE HEX(your_column) LIKE '%96%';
    
    Check if the custom operator correctly converts this byte to the proper en dash character.
  2. Run a single task: Trigger just the extract task for the problematic table to confirm it no longer throws the UnicodeDecodeError.
  3. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 08:00:51