airflow.providers.google.cloud.hooks.gcs

此模块包含一个 Google Cloud Storage Hook。

属性

RT

T

FParams

List

DEFAULT_TIMEOUT

PROVIDE_BUCKET

GCSHook

使用 Google Cloud 连接与 Google Cloud Storage 交互。

GCSAsyncHook

GCSAsyncHook 在触发器工作进程上运行,继承自 GoogleBaseAsyncHook。

函数

gcs_object_is_directory(bucket)

如果给定的 Google Cloud Storage URL(gs://<bucket>/<blob>)是目录或空存储桶,则返回 True。

parse_json_from_gcs(gcp_conn_id, file_uri[, ...])

从 Google Cloud Storage 下载并解析 JSON 文件。

模块内容

airflow.providers.google.cloud.hooks.gcs.RT[source]
airflow.providers.google.cloud.hooks.gcs.T[source]
airflow.providers.google.cloud.hooks.gcs.FParams[source]
airflow.providers.google.cloud.hooks.gcs.List[source]
airflow.providers.google.cloud.hooks.gcs.DEFAULT_TIMEOUT = 60[source]
airflow.providers.google.cloud.hooks.gcs.PROVIDE_BUCKET: str = None[source]
class airflow.providers.google.cloud.hooks.gcs.GCSHook(gcp_conn_id='google_cloud_default', impersonation_chain=None, **kwargs)[source]

Bases: airflow.providers.google.common.hooks.base_google.GoogleBaseHook

使用 Google Cloud 连接与 Google Cloud Storage 交互。

get_conn()[source]

返回一个 Google Cloud Storage 服务对象。

copy(source_bucket, source_object, destination_bucket=None, destination_object=None)[source]

将对象从一个存储桶复制到另一个存储桶,可在需要时进行重命名。

destination_bucket 或 destination_object 可以省略,此时使用源存储桶/对象,但两者不能同时省略。

参数:
  • source_bucket (str) – 要复制对象的源存储桶。

  • source_object (str) – 要复制的对象。

  • destination_bucket (str | None) – 复制目标的存储桶。若省略,则使用相同的存储桶。

  • destination_object (str | None) – 若提供,则为目标对象的(重命名后)路径。若省略,则使用相同的名称。

rewrite(source_bucket, source_object, destination_bucket, destination_object=None)[source]

功能类似 copy;支持超过 5 TB 的文件以及跨区域或跨存储类的复制。

destination_object 可省略,此时使用 source_object。

参数:
  • source_bucket (str) – 要复制对象的源存储桶。

  • source_object (str) – 要复制的对象。

  • destination_bucket (str) – 复制目标的存储桶。

  • destination_object (str | None) – 若提供,则为目标对象的(重命名后)路径。若省略,则使用相同的名称。

download(bucket_name: str, object_name: str, filename: None = None, chunk_size: int | None = None, timeout: int | None = DEFAULT_TIMEOUT, num_max_attempts: int | None = 1, user_project: str | None = None) bytes[source]
download(bucket_name: str, object_name: str, filename: str, chunk_size: int | None = None, timeout: int | None = DEFAULT_TIMEOUT, num_max_attempts: int | None = 1, user_project: str | None = None) str

从 Google Cloud Storage 下载文件。

如果未提供 filename,运算符会将文件加载到内存并返回其内容;如果提供了 filename,则会将文件写入指定位置并返回该路径。对于超过可用内存的文件,建议写入磁盘。

参数:
  • bucket_name – 要获取的存储桶。

  • object_name – 要获取的对象。

  • filename – 若设置,则为文件应写入的本地路径。

  • chunk_size – Blob 的块大小。

  • timeout – 请求超时时间(秒)。

  • num_max_attempts – 下载文件的尝试次数。

  • user_project – 用于计费的 Google Cloud 项目标识符。Requester Pays 存储桶必需。

download_as_byte_array(bucket_name, object_name, chunk_size=None, timeout=DEFAULT_TIMEOUT, num_max_attempts=1)[source]

从 Google Cloud Storage 下载文件。

如果未提供 filename,运算符会将文件加载到内存并返回其内容;如果提供了 filename,则会将文件写入指定位置并返回该路径。对于超过可用内存的文件,建议写入磁盘。

参数:
  • bucket_name (str) – 要获取的存储桶。

  • object_name (str) – 要获取的对象。

  • chunk_size (int | None) – Blob 的块大小。

  • timeout (int | None) – 请求超时时间(秒)。

  • num_max_attempts (int | None) – 下载文件的尝试次数。

provide_file(bucket_name=PROVIDE_BUCKET, object_name=None, object_url=None, dir=None, user_project=None)[source]

将文件下载到临时目录并返回文件句柄。

可以通过传入 bucket_name 与 object_name 参数,或仅传入 object_url 参数来使用此方法。

参数:
  • bucket_name (str) – 要获取的存储桶。

  • object_name (str | None) – 要获取的对象。

  • object_url (str | None) – 文件引用 URL,必须以“gs://”开头。

  • dir (str | None) – 下载文件的临时子目录(传递给 NamedTemporaryFile)。

  • user_project (str | None) – 用于计费的 Google Cloud 项目标识符。Requester Pays 存储桶必需。

返回:

文件句柄

返回类型:

collections.abc.Generator[IO[bytes], None, None]

provide_file_and_upload(bucket_name=PROVIDE_BUCKET, object_name=None, object_url=None, user_project=None)[source]

创建临时文件,返回文件句柄并在关闭时自动上传文件内容。

可以通过传入 bucket_name 与 object_name 参数,或仅传入 object_url 参数来使用此方法。

参数:
  • bucket_name (str) – 要获取的存储桶。

  • object_name (str | None) – 要获取的对象。

  • object_url (str | None) – 文件引用 URL,必须以“gs://”开头。

  • user_project (str | None) – 用于计费的 Google Cloud 项目标识符。Requester Pays 存储桶必需。

返回:

文件句柄

返回类型:

collections.abc.Generator[IO[bytes], None, None]

upload(bucket_name, object_name, filename=None, data=None, mime_type=None, gzip=False, encoding='utf-8', chunk_size=None, timeout=DEFAULT_TIMEOUT, num_max_attempts=1, metadata=None, cache_control=None, user_project=None)[source]

将本地文件或字符串/字节数据上传到 Google Cloud Storage。

参数:
  • bucket_name (str) – 要上传到的存储桶。

  • object_name (str) – 上传文件时设置的对象名称。

  • filename (str | None) – 要上传的本地文件路径。

  • data (str | bytes | None) – 要上传的文件数据(字符串或字节)。

  • mime_type (str | None) – 上传文件时设置的 MIME 类型。

  • gzip (bool) – 是否在上传时压缩本地文件或文件数据

  • encoding (str) – 当 data 为字符串时的字节编码。

  • chunk_size (int | None) – Blob 的块大小。

  • timeout (int | None) – 请求超时时间(秒)。

  • num_max_attempts (int) – 上传文件的最大尝试次数。

  • metadata (dict | None) – 与文件一起上传的元数据。

  • cache_control (str | None) – Cache-Control 元数据字段。

  • user_project (str | None) – 用于计费的 Google Cloud 项目标识符。Requester Pays 存储桶必需。

exists(bucket_name, object_name, retry=DEFAULT_RETRY, user_project=None)[source]

检查 Google Cloud Storage 中文件是否存在。

参数:
  • bucket_name (str) – 对象所在的 Google Cloud Storage 存储桶。

  • object_name (str) – 要检查的 blob 名称。

  • retry (google.api_core.retry.Retry) – (可选)RPC 的重试策略。

  • user_project (str | None) – 用于计费的 Google Cloud 项目标识符。Requester Pays 存储桶必需。

get_blob_update_time(bucket_name, object_name)[source]

获取 Google Cloud Storage 中文件的更新时间。

参数:
  • bucket_name (str) – 对象所在的 Google Cloud Storage 存储桶。

  • object_name (str) – 要获取更新时间的 blob 名称。

is_updated_after(bucket_name, object_name, ts)[source]

检查 blob 是否在给定时间之后更新。

参数:
  • bucket_name (str) – 对象所在的 Google Cloud Storage 存储桶。

  • object_name (str) – 要检查的对象名称。

  • ts (datetime.datetime) – 用于比较的时间戳。

is_updated_between(bucket_name, object_name, min_ts, max_ts)[source]

检查 blob 是否在给定时间之后更新。

参数:
  • bucket_name (str) – 对象所在的 Google Cloud Storage 存储桶。

  • object_name (str) – 要检查的对象名称。

  • min_ts (datetime.datetime) – 最小时间戳。

  • max_ts (datetime.datetime) – 最大时间戳。

is_updated_before(bucket_name, object_name, ts)[source]

检查 blob 是否在给定时间之前更新。

参数:
  • bucket_name (str) – 对象所在的 Google Cloud Storage 存储桶。

  • object_name (str) – 要检查的对象名称。

  • ts (datetime.datetime) – 用于比较的时间戳。

is_older_than(bucket_name, object_name, seconds)[source]

检查对象是否比指定时间更旧。

参数:
  • bucket_name (str) – 对象所在的 Google Cloud Storage 存储桶。

  • object_name (str) – 要检查的对象名称。

  • seconds (int) – 用于比较的秒数。

delete(bucket_name, object_name, ignore_error=False)[source]

从存储桶中删除对象。

参数:
  • bucket_name (str) – 对象所在的存储桶名称。

  • object_name (str) – 要删除的对象名称。

  • ignore_error (bool) – (可选)是否忽略 NotFound 异常,默认 False。

get_bucket(bucket_name)[source]

从 Google Cloud Storage 获取存储桶对象。

参数:

bucket_name (str) – 存储桶名称。

delete_bucket(bucket_name, force=False, user_project=None)[source]

删除 Google Cloud Storage 中的存储桶对象。

参数:
  • bucket_name (str) – 要删除的存储桶名称。

  • force (bool) – False 时不允许删除非空存储桶,设为 True 可强制删除非空存储桶。

  • user_project (str | None) – 用于计费的 Google Cloud 项目标识符。Requester Pays 存储桶必需。

list(bucket_name, versions=None, max_results=None, prefix=None, delimiter=None, match_glob=None, user_project=None)[source]

列出存储桶中所有对象,可使用单个或多个前缀进行筛选。

参数:
  • bucket_name (str) – 存储桶名称。

  • versions (bool | None) – 为 True 时列出对象的所有版本。

  • max_results (int | None) – 单页响应返回的最大条目数。

  • prefix (str | List[str] | None) – 字符串或字符串列表,用于过滤名称以其/它们开头的对象

  • delimiter (str | None) – (已弃用)根据分隔符过滤对象(例如 ‘.csv’)

  • match_glob (str | None) – (可选)根据字符串提供的 glob 模式过滤对象(例如 '**/*/.json')。

  • user_project (str | None) – 用于计费的 Google Cloud 项目标识符。Requester Pays 存储桶必需。

返回:

匹配过滤条件的对象名称流

list_by_timespan(bucket_name, timespan_start, timespan_end, versions=None, max_results=None, prefix=None, delimiter=None, match_glob=None)[source]

列出存储桶中具有给定字符串前缀且在时间范围内更新的所有对象。

参数:
  • bucket_name (str) – 存储桶名称。

  • timespan_start (datetime.datetime) – 返回在此日期时间(UTC)或之后更新的对象

  • timespan_end (datetime.datetime) – 返回在此日期时间(UTC)之前更新的对象

  • versions (bool | None) – 为 True 时列出对象的所有版本。

  • max_results (int | None) – 单页响应返回的最大条目数。

  • prefix (str | None) – 用于过滤名称以此前缀开头的对象的前缀字符串

  • delimiter (str | None) – (已弃用)根据分隔符过滤对象(例如 ‘.csv’)

  • match_glob (str | None) – (可选)根据字符串提供的 glob 模式过滤对象(例如 '**/*/.json')。

返回:

匹配过滤条件的对象名称流

返回类型:

List[str]

get_size(bucket_name, object_name)[source]

获取 Google Cloud Storage 中文件的大小。

参数:
  • bucket_name (str) – 包含 blob_name 的 Google Cloud Storage 存储桶。

  • object_name (str) – 要在 Google Cloud Storage bucket_name 中检查的对象名称。

get_crc32c(bucket_name, object_name)[source]

获取 Google Cloud Storage 中对象的 CRC32c 校验和。

参数:
  • bucket_name (str) – 包含 blob_name 的 Google Cloud Storage 存储桶。

  • object_name (str) – 要在 Google Cloud Storage bucket_name 中检查的对象名称。

get_md5hash(bucket_name, object_name)[source]

获取 Google Cloud Storage 中对象的 MD5 哈希值。

参数:
  • bucket_name (str) – 包含 blob_name 的 Google Cloud Storage 存储桶。

  • object_name (str) – 要在 Google Cloud Storage bucket_name 中检查的对象名称。

get_metadata(bucket_name, object_name)[source]

获取 Google Cloud Storage 中对象的元数据。

参数:
  • bucket_name (str) – 对象所在的 Google Cloud Storage 存储桶名称。

  • object_name (str) – 包含所需元数据的对象名称

返回:

与对象关联的元数据

返回类型:

dict | None

create_bucket(bucket_name, resource=None, storage_class='MULTI_REGIONAL', location='US', project_id=PROVIDE_PROJECT_ID, labels=None)[source]

创建一个新存储桶。

Google Cloud Storage 使用平面命名空间,因此不能创建已被使用的存储桶名称。

另请参阅

有关更多信息,请参阅存储桶命名指南: https://cloud.google.com/storage/docs/bucketnaming.html#requirements

参数:
  • bucket_name (str) – 存储桶的名称。

  • resource (dict | None) – 用于创建存储桶的可选字典参数。有关可用参数的信息,请参阅 Cloud Storage API 文档: https://cloud.google.com/storage/docs/json_api/v1/buckets/insert

  • storage_class (str) – —

    这定义了存储桶中对象的存储方式,并决定 SLA 和存储费用。取值包括

    • MULTI_REGIONAL

    • REGIONAL

    • STANDARD

    • NEARLINE

    • COLDLINE.

    如果在创建存储桶时未指定此值,默认使用 STANDARD。

  • location (str) –

    存储桶的地点。对象数据存储在该地区的实际存储中。默认值为 US。

  • project_id (str) – Google Cloud 项目的 ID。

  • labels (dict | None) – 用户提供的标签,以键/值对形式。

返回:

如果成功,返回存储桶的 id

返回类型:

str

insert_bucket_acl(bucket_name, entity, role, user_project=None)[source]

在指定的 bucket_name 上创建新的 ACL 条目。

See: https://cloud.google.com/storage/docs/json_api/v1/bucketAccessControls/insert

参数:
  • bucket_name (str) – 存储桶的名称。

  • entity (str) – 拥有权限的实体,形式之一:user-userId、user-email、group-groupId、group-email、domain-domain、project-team-projectId、allUsers、allAuthenticatedUsers。参见: https://cloud.google.com/storage/docs/access-control/lists#scopes

  • role (str) – 实体的访问权限。可接受的值有:“OWNER”、 “READER”、 “WRITER”。

  • user_project (str | None) – (可选)为此请求计费的项目。对 Requester Pays 存储桶是必需的。

insert_object_acl(bucket_name, object_name, entity, role, generation=None, user_project=None)[source]

在指定对象上创建新的 ACL 条目。

See: https://cloud.google.com/storage/docs/json_api/v1/objectAccessControls/insert

参数:
  • bucket_name (str) – 存储桶的名称。

  • object_name (str) – 对象的名称。有关如何对对象名称进行 URL 编码以确保路径安全的信息,请参见: https://cloud.google.com/storage/docs/json_api/#encoding

  • entity (str) – 拥有权限的实体,形式之一:user-userId、user-email、group-groupId、group-email、domain-domain、project-team-projectId、allUsers、allAuthenticatedUsers。参见: https://cloud.google.com/storage/docs/access-control/lists#scopes

  • role (str) – 实体的访问权限。可接受的值有:“OWNER”、 “READER”。

  • generation (int | None) – 可选。如果提供,则选择此对象的特定修订版本。

  • user_project (str | None) – (可选)为此请求计费的项目。对 Requester Pays 存储桶是必需的。

compose(bucket_name, source_objects, destination_object)[source]

在同一存储桶中将现有对象列表合成为一个新对象。

当前一次操作最多支持合并 32 个对象。

https://cloud.google.com/storage/docs/json_api/v1/objects/compose

参数:
  • bucket_name (str) – 包含源对象的存储桶名称。该存储桶也用于存储合成后的目标对象。

  • source_objects (List[str]) – 将要合成为单个对象的源对象列表。

  • destination_object (str) – 目标对象的路径(如果提供)。

sync_to_local_dir(bucket_name, local_dir, prefix=None, delete_stale=False)[source]

将文件从 GCS 存储桶下载到本地目录。

它会下载给定 prefix 下的所有文件,并在 local_dir 中创建相应的目录结构。

如果 delete_staleTrue,则会删除 GCS 存储桶中不存在的本地文件。

参数:
  • bucket_name (str) – GCS 存储桶的名称。

  • local_dir (str | pathlib.Path) – 文件将下载到的本地目录。

  • prefix (str | None) – 要下载的文件的前缀。

  • delete_stale (bool) – 如果为 True,则删除在存储桶中不存在的本地文件。

sync(source_bucket, destination_bucket, source_object=None, destination_object=None, recursive=True, allow_overwrite=False, delete_extra_files=False)[source]

同步存储桶的内容。

参数 source_objectdestination_object 描述根同步目录。如果未传入,则同步整个存储桶;如果传入,则应指向目录。

注意

不支持单个文件的同步。仅支持整个目录的同步。

参数:
  • source_bucket (str) – 包含源对象的存储桶名称。

  • destination_bucket (str) – 包含目标对象的存储桶名称。

  • source_object (str | None) – 源存储桶中的根同步目录。

  • destination_object (str | None) – 目标存储桶中的根同步目录。

  • recursive (bool) – 如果为 True,则会考虑子目录

  • recursive – 如果为 True,则会考虑子目录

  • allow_overwrite (bool) – 如果为 True,则在发现不匹配的文件时覆盖文件。默认不允许覆盖文件。

  • delete_extra_files (bool) –

    如果为 True,则删除目标中不存在的源额外文件。默认不删除额外文件。

    注意

    如果指定错误的源/目标组合,此选项可能会快速删除数据。

返回:

返回类型:

airflow.providers.google.cloud.hooks.gcs.gcs_object_is_directory(bucket)[source]

如果给定的 Google Cloud Storage URL(gs://<bucket>/<blob>)是目录或空存储桶,则返回 True。

airflow.providers.google.cloud.hooks.gcs.parse_json_from_gcs(gcp_conn_id, file_uri, impersonation_chain=None)[source]

从 Google Cloud Storage 下载并解析 JSON 文件。

参数:
  • gcp_conn_id (str) – Airflow Google Cloud 连接 ID。

  • file_uri (str) – JSON 文件的完整路径,例如:gs://test-bucket/dir1/dir2/file

class airflow.providers.google.cloud.hooks.gcs.GCSAsyncHook(**kwargs)[source]

基础类: airflow.providers.google.common.hooks.base_google.GoogleBaseAsyncHook

GCSAsyncHook 在触发器工作进程上运行,继承自 GoogleBaseAsyncHook。

sync_hook_class[source]
async get_storage_client(session)[source]

返回一个 Google Cloud Storage 服务对象。

此条目是否有帮助?