S3 插件¶
S3 插件(datus-s3-plugin)让 Datus agent 能够浏览、查询和搬运对象存储数据。
它既支持 Amazon S3,也支持 MinIO、阿里云 OSS 等 S3 兼容存储;当工作流需要读写
s3:// URI 时(比如部署 Airflow DAG),其他插件都依赖它承担传输。
安装¶
要求 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:ListBucket 和 s3: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 部署的传输层¶
Airflow 和 MWAA 插件刻意不含对象存储客户端。
只要 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 的发现和加载方式