Google Cloud Storage 操作符¶
Cloud Storage 允许在全球范围内随时存储和检索任意数量的数据。您可以将 Cloud Storage 用于多种场景,包括提供网站内容、存储用于归档和灾难恢复的数据,或通过直接下载向用户分发大型数据对象。
请参阅 Google Transfer 操作符 以获取针对 Google Cloud Storage 的专用传输操作符列表。
前置任务¶
要使用这些操作符,您需要完成以下几项工作
使用 云控制台 选择或创建一个云平台项目。
为您的项目启用计费,方法请参阅 Google Cloud 文档。
启用 API,方法请参阅 云控制台文档。
通过 pip 安装 API 库。
pip install 'apache-airflow[google]'有关详细信息,请参阅 安装。
操作符¶
GCSTimeSpanFileTransformOperator¶
使用 GCSTimeSpanFileTransformOperator 来转换在特定时间段(数据区间)内被修改的文件。时间段由开始和结束时间戳定义。如果 Dag 没有调度下一个 Dag 实例,则时间段结束为无限,这意味着操作符会处理所有早于 data_interval_start 的文件。
tests/system/google/cloud/gcs/example_gcs_transform_timespan.py
gcs_timespan_transform_files_task = GCSTimeSpanFileTransformOperator(
task_id="gcs_timespan_transform_files",
source_bucket=BUCKET_NAME_SRC,
source_prefix=SOURCE_PREFIX,
source_gcp_conn_id=SOURCE_GCP_CONN_ID,
destination_bucket=BUCKET_NAME_DST,
destination_prefix=DESTINATION_PREFIX,
destination_gcp_conn_id=DESTINATION_GCP_CONN_ID,
transform_script=["python", TRANSFORM_SCRIPT_PATH],
)
gcs_timespan_transform_files_parallel = GCSTimeSpanFileTransformOperator(
task_id="gcs_timespan_transform_files_parallel",
source_bucket=BUCKET_NAME_SRC,
source_prefix=SOURCE_PREFIX,
source_gcp_conn_id=SOURCE_GCP_CONN_ID,
destination_bucket=BUCKET_NAME_DST,
destination_prefix=DESTINATION_PREFIX_PARALLEL,
destination_gcp_conn_id=DESTINATION_GCP_CONN_ID,
transform_script=["python", TRANSFORM_SCRIPT_PATH],
max_download_workers=2,
max_upload_workers=2,
)
GCSBucketCreateAclEntryOperator¶
在指定的存储桶上创建新的 ACL 条目。
有关参数定义,请查看 GCSBucketCreateAclEntryOperator
使用该运算符¶
gcs_bucket_create_acl_entry_task = GCSBucketCreateAclEntryOperator(
bucket=BUCKET_NAME,
entity=GCS_ACL_ENTITY,
role=GCS_ACL_BUCKET_ROLE,
task_id="gcs_bucket_create_acl_entry_task",
)
模板化¶
template_fields: Sequence[str] = (
"bucket",
"entity",
"role",
"user_project",
"gcp_conn_id",
"impersonation_chain",
)
更多信息¶
请参阅 Google Cloud Storage 文档 以在存储桶中创建新的 ACL 条目。
GCSObjectCreateAclEntryOperator¶
在指定的对象上创建新的 ACL 条目。
有关参数定义,请查看 GCSObjectCreateAclEntryOperator
使用该操作符¶
gcs_object_create_acl_entry_task = GCSObjectCreateAclEntryOperator(
bucket=BUCKET_NAME,
object_name=FILE_NAME,
entity=GCS_ACL_ENTITY,
role=GCS_ACL_OBJECT_ROLE,
task_id="gcs_object_create_acl_entry_task",
)
模板化¶
template_fields: Sequence[str] = (
"bucket",
"object_name",
"entity",
"generation",
"role",
"user_project",
"gcp_conn_id",
"impersonation_chain",
)
更多信息¶
请参阅 Google Cloud Storage 插入文档,以在 ObjectAccess 中创建 ACL 条目。
GCSListObjectsOperator¶
使用 GCSListObjectsOperator 列出 Google Cloud Storage 存储桶中的对象。您可以选择指定前缀,仅列出名称以该前缀开头的对象,亦可指定分隔符以模拟目录结构。
tests/system/google/cloud/gcs/example_gcs_copy_delete.py
list_buckets = GCSListObjectsOperator(task_id="list_buckets", bucket=BUCKET_NAME_SRC)
GCSDeleteObjectsOperator¶
使用 GCSDeleteObjectsOperator 从 Google Cloud Storage 存储桶中删除一个或多个对象。
tests/system/google/cloud/gcs/example_gcs_copy_delete.py
delete_files = GCSDeleteObjectsOperator(
task_id="delete_files", bucket_name=BUCKET_NAME_SRC, objects=[FILE_NAME]
)
Deleting Bucket
Deleting Bucket 允许您从 Google Cloud Storage 中移除存储桶对象。此操作通过 GCSDeleteBucketOperator 操作符执行。
tests/system/google/cloud/gcs/example_gcs_upload_download.py
delete_bucket = GCSDeleteBucketOperator(task_id="delete_bucket", bucket_name=BUCKET_NAME)
您可以在 bucket_name、gcp_conn_id、impersonation_chain、user_project 参数上使用 Jinja 模板,从而动态确定这些值。
参考¶
欲了解更多信息,请参阅
传感器¶
GCSObjectExistenceSensor¶
使用 GCSObjectExistenceSensor 等待(轮询) Google Cloud Storage 中文件的存在。
gcs_object_exists = GCSObjectExistenceSensor(
bucket=DESTINATION_BUCKET_NAME,
object=FILE_NAME,
task_id="gcs_object_exists_task",
)
如果您希望在传感器运行期间释放工作节点槽位,也可以在此操作符中使用可延迟(deferrable)模式。
gcs_object_exists_defered = GCSObjectExistenceSensor(
bucket=DESTINATION_BUCKET_NAME, object=FILE_NAME, task_id="gcs_object_exists_defered", deferrable=True
)
GCSObjectsWithPrefixExistenceSensor¶
使用 GCSObjectsWithPrefixExistenceSensor 等待(轮询) Google Cloud Storage 中具有指定前缀的文件的存在。
gcs_object_with_prefix_exists = GCSObjectsWithPrefixExistenceSensor(
bucket=DESTINATION_BUCKET_NAME,
prefix=FILE_NAME[:5],
task_id="gcs_object_with_prefix_exists_task",
)
如果您希望此传感器异步运行,从而更高效地利用 Airflow 部署中的资源,可将 deferrable 参数设置为 True。不过,需要启用 triggerer 组件才能使此功能生效。
gcs_object_with_prefix_exists_async = GCSObjectsWithPrefixExistenceSensor(
bucket=DESTINATION_BUCKET_NAME,
prefix=FILE_NAME[:5],
task_id="gcs_object_with_prefix_exists_task_async",
deferrable=True,
)
GCSUploadSessionCompleteSensor¶
使用 GCSUploadSessionCompleteSensor 检查 Google Cloud Storage 中具有指定前缀的文件数量是否有变化。
gcs_upload_session_complete = GCSUploadSessionCompleteSensor(
bucket=DESTINATION_BUCKET_NAME,
prefix=FILE_NAME,
inactivity_period=15,
min_objects=1,
allow_delete=True,
previous_objects=set(),
task_id="gcs_upload_session_complete_task",
)
如果您希望在传感器运行期间释放工作节点槽位,可将参数 deferrable 设置为 True。
gcs_upload_session_async_complete = GCSUploadSessionCompleteSensor(
bucket=DESTINATION_BUCKET_NAME,
prefix=FILE_NAME,
inactivity_period=15,
min_objects=1,
allow_delete=True,
previous_objects=set(),
task_id="gcs_upload_session_async_complete",
deferrable=True,
)
GCSObjectUpdateSensor¶
使用 GCSObjectUpdateSensor 检查 Google Cloud Storage 中的对象是否已更新。
gcs_update_object_exists = GCSObjectUpdateSensor(
bucket=DESTINATION_BUCKET_NAME,
object=FILE_NAME,
task_id="gcs_object_update_sensor_task",
)
如果您希望此传感器异步运行,以更高效利用 Airflow 部署中的资源,可将 deferrable 参数设为 True。不过,需要启用 triggerer 组件才能使该功能正常工作。
gcs_update_object_exists_async = GCSObjectUpdateSensor(
bucket=DESTINATION_BUCKET_NAME,
object=FILE_NAME,
task_id="gcs_object_update_sensor_task_async",
deferrable=True,
)
更多信息¶
传感器有不同的模式来决定任务执行期间资源的行为。请参阅 Airflow 传感器文档,了解使用传感器的最佳实践。
参考¶
欲了解更多信息,请参阅