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
相关产品推荐
相关产品推荐

