跳转至

S3 插件

S3 插件(datus-s3-plugin)让 Datus agent 能够浏览、查询和搬运对象存储数据。 它既支持 Amazon S3,也支持 MinIO、阿里云 OSS 等 S3 兼容存储;当工作流需要读写 s3:// URI 时(比如部署 Airflow DAG),其他插件都依赖它承担传输。

安装

datus plugin install datus-s3-plugin

要求 datus-agent >= 0.3.8。其他安装来源和 profile 管理方式见 插件

Skills

Skill 作用
s3 浏览、读取、查询和搬运对象
s3-setup 创建 profile(区域、凭据、默认桶、兼容模式)

s3

借助这个 skill,agent 可以:

  • 浏览和读取 —— 列出桶和前缀、查看对象元数据、打印对象内容或前几行、生成 预签名 URL(预签名的 PUT URL 等于写权限,会被当作凭据对待);
  • 就地查询 —— 用 S3 Select SQL 直接查询单个 CSV、JSON 或 Parquet 对象, 不需要下载,比如直接在桶上按分组数行数;
  • 搬运数据 —— 上传下载文件、同步目录(只传输新增和变化的文件)、移动对象、 删除对象;
  • 安全发布产物 —— skill 内置了上传规范:上传后读回元数据做验证,每个构建 发布到带版本号的 key 而不是覆盖可变 key,绝不把凭据留在上传的文件里。

Profile 设置了 kms_key_id 时,写入使用 SSE-KMS 加密。aliyun-oss 兼容模式下, S3 Select 和 SSE-KMS 会被明确拒绝。

注意:AWS 已不再向新客户开放 S3 Select,只有 2024 年 7 月 25 日之前用过它的 账号可以继续调用。新账号建议改用 Amazon Athena 查询,或下载对象后在本地过滤; MinIO 仍然支持 S3 Select API。

s3-setup

让 agent 帮你配置插件,这个 skill 会收集区域、凭据来源(默认 AWS 链、命名 profile、以 ${ENV_VAR} 引用的密钥,或要 assume 的角色)、供裸 key 使用的可选 默认桶,以及可选的 SSE-KMS 密钥,然后用一次只读列表验证 profile。MinIO 只需要 endpoint_url;阿里云 OSS 使用专门的兼容字段。生成的 profile 形如:

agent:
  plugins:
    s3:
      prod:
        default: true
        region: us-east-1
        bucket: my-data-lake  # 可选的默认桶
        # 凭据:标准 AWS 链,或 profile / 密钥 / role_arn
      minio-local:
        endpoint_url: http://minio:9000
        access_key_id: ${MINIO_ACCESS_KEY}
        secret_access_key: ${MINIO_SECRET_KEY}

IAM 主体至少需要 s3:ListBuckets3:GetObject;列出账号下的所有桶还需要 s3:ListAllMyBuckets(其 Resource 必须是 "*")。上传需要 s3:PutObject, 删除需要 s3:DeleteObject

在 agent 中使用

配置好 profile 后,可以这样提出请求:

  • 「s3://my-lake/events/ 里有什么?」 —— 列出该前缀下的对象。
  • 「把 s3://my-lake/events/day.csv 的前十行给我看看。」
  • 「这个文件里每个 region 各有多少条记录?」 —— 一次 S3 Select 查询, 什么都不用下载。
  • 「把 ./artifacts 同步到 s3://my-lake/artifacts/v1/ 并验证上传结果。」
  • 「帮我把 s3 插件配到我们的 MinIO 上。」 —— 触发 s3-setup

权限系统下,读取、查询和预签名不需要确认; 上传、同步和移动在 normal 模式下确认一次;删除任何时候都要确认。

编排工作流

DAG 部署的传输层

AirflowMWAA 插件刻意不含对象存储客户端。 只要 DAG 的部署、导出、备份或迁移涉及 s3:// URI,agent 就会把传输交给本插件:

  • 部署 —— 「把 dags/sales_daily.py 部署到 prod Airflow」:agent 通过本插件 把文件上传到调度器的 dags_folder(或 MWAA 环境的 DAG 前缀),再通过调度器 侧的插件验证 DAG 解析成功。
  • 备份 —— 「把 prod 的所有活跃 DAG 导出到 s3://backup/airflow/」:Airflow 或 MWAA 的导出 skill 通过调度器的 API 收集源码,把上传交给本插件。

凭据彼此隔离:本插件持有存储凭据,调度器插件持有 Airflow 凭据,双方都看不到 对方的。

发布流水线产物

同样的模式适用于流水线从对象存储消费的一切东西——作业 JAR、渲染好的配置、导出 的查询结果。让 agent 把目录发布到带版本号的前缀下:它只同步有变化的文件,然后 读回上传结果做验证。

相关文档

  • 插件 —— 安装来源、profile、项目启用与权限
  • Airflow 插件 —— 自建 Airflow 上的 DAG 部署
  • MWAA 插件 —— Amazon MWAA 上的 DAG 部署
  • Skills —— skills 的发现和加载方式