能否在Airflow的GoogleCloudStorageToBigQueryOperator中使用通配符?
问题
我在GCS的某个文件夹下存放了一系列文件,文件名如下:
file_sample_1.json file_sample_2.json file_sample_3.json ... file_sample_n.json
希望通过Airflow的GoogleCloudStorageToBigQueryOperator将这些文件导入BigQuery。以下是我的代码:
def create_operator_write_init(): return GoogleCloudStorageToBigQueryOperator( task_id = 'test_ingest_to_bq', bucket = 'sample-bucket-dev-202211', source_objects = 'file_sample_1.json', destination_project_dataset_table = 'sample_destination_table', create_disposition = "CREATE_IF_NEEDED", write_disposition = "WRITE_TRUNCATE", source_format = "NEWLINE_DELIMITED_JSON", schema_fields = [ {"name": "id", "type": "INTEGER", "mode": "NULLABLE"}, {"name": "created_at", "type": "TIMESTAMP", "mode": "NULLABLE"}, {"name": "updated_at", "type": "TIMESTAMP", "mode": "NULLABLE"}, ] )
当前代码可正常导入单个文件,但我需要通过通配符批量指定源文件,请问能否使用file_sample_*.json这类格式实现?
回答
完全可以使用file_sample_*.json这类通配符格式批量指定GCS源文件,只需修改source_objects参数的值即可。
修改后的代码如下:
def create_operator_write_init(): return GoogleCloudStorageToBigQueryOperator( task_id = 'test_ingest_to_bq', bucket = 'sample-bucket-dev-202211', source_objects = 'file_sample_*.json', # 替换为通配符格式 destination_project_dataset_table = 'sample_destination_table', create_disposition = "CREATE_IF_NEEDED", write_disposition = "WRITE_TRUNCATE", source_format = "NEWLINE_DELIMITED_JSON", schema_fields = [ {"name": "id", "type": "INTEGER", "mode": "NULLABLE"}, {"name": "created_at", "type": "TIMESTAMP", "mode": "NULLABLE"}, {"name": "updated_at", "type": "TIMESTAMP", "mode": "NULLABLE"}, ] )
补充说明
- 若文件存放在GCS子文件夹中,比如
data/file_sample_*.json,直接在source_objects里指定带路径的通配符即可。 - 该操作符支持标准GCS通配符规则:
*匹配任意长度任意字符,?匹配单个任意字符,也可使用[abc]这类字符集匹配。 - 需确保所有匹配到的文件格式统一,符合
NEWLINE_DELIMITED_JSON要求,避免导入失败。
内容的提问来源于stack exchange,提问作者Iqbal
相关产品推荐
相关产品推荐

