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

GCP BigQuery导出CSV如何指定CRLF换行?含Airflow实现方案

BigQuery导出CSV指定CRLF换行符及Airflow实现方案

核心结论

BigQuery的EXPORT DATA语句没有直接指定换行符的参数,默认导出的CSV文件使用LF作为换行符,无法通过语句内的options配置改为CRLF,需要通过后续处理实现。

Airflow实现方案

方案一:Python Operator实现

先通过BigQuery导出算子将数据以LF格式导出到GCS,再用Python代码读取文件、替换换行符为CRLF后重新上传。

示例代码:

from airflow import DAG
from airflow.providers.google.cloud.operators.bigquery import BigQueryExportOperator
from airflow.operators.python import PythonOperator
from google.cloud import storage
from datetime import datetime

def convert_lf_to_crlf(bucket_name, file_path):
    storage_client = storage.Client()
    bucket = storage_client.bucket(bucket_name)
    blob = bucket.blob(file_path)
    
    # 流式处理大文件,避免内存溢出
    with blob.open("r") as f_in:
        with blob.open("w") as f_out:
            for line in f_in:
                # 替换LF为CRLF
                f_out.write(line.rstrip('\n') + '\r\n')

with DAG(
    dag_id="bq_export_crlf",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    # 第一步:导出数据到GCS(LF格式)
    export_task = BigQueryExportOperator(
        task_id="export_bq_to_gcs",
        source_project_dataset_table="<your-project>.<your-dataset>.<your-table>",
        destination_cloud_storage_uris=["gs://<bucket>/<file_name>.txt"],
        export_format="CSV",
        field_delimiter="~",
        print_header=False,
        overwrite=True
    )

    # 第二步:转换换行符为CRLF
    convert_task = PythonOperator(
        task_id="convert_lf_to_crlf",
        python_callable=convert_lf_to_crlf,
        op_kwargs={
            "bucket_name": "<bucket>",
            "file_path": "<file_name>.txt"
        }
    )

    export_task >> convert_task

方案二:Bash Operator实现

利用gsutil下载GCS文件到Airflow worker本地,通过命令行工具转换换行符后再上传回GCS。

示例代码:

from airflow import DAG
from airflow.providers.google.cloud.operators.bigquery import BigQueryExportOperator
from airflow.operators.bash import BashOperator
from datetime import datetime

with DAG(
    dag_id="bq_export_crlf_bash",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    export_task = BigQueryExportOperator(
        task_id="export_bq_to_gcs",
        source_project_dataset_table="<your-project>.<your-dataset>.<your-table>",
        destination_cloud_storage_uris=["gs://<bucket>/<file_name>.txt"],
        export_format="CSV",
        field_delimiter="~",
        print_header=False,
        overwrite=True
    )

    convert_task = BashOperator(
        task_id="convert_lf_to_crlf",
        bash_command="""
            # 下载文件到本地
            gsutil cp gs://<bucket>/<file_name>.txt /tmp/temp_file.txt
            # 用sed替换每行末尾添加\r,实现LF转CRLF
            sed 's/$/\r/' /tmp/temp_file.txt > /tmp/converted_file.txt
            # 上传回GCS覆盖原文件
            gsutil cp /tmp/converted_file.txt gs://<bucket>/<file_name>.txt
            # 清理本地临时文件
            rm /tmp/temp_file.txt /tmp/converted_file.txt
        """
    )

    export_task >> convert_task

注意:若Airflow worker环境未安装sed,可改用unix2dos命令(unix2dos /tmp/temp_file.txt),需确保环境已预装该工具。

内容的提问来源于stack exchange,提问作者Mani Shankar.S

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 00:02:54