| name | create-pipeline-task |
| description | 通过 CLI 创建集成管道任务(数据同步/数据搬运/ETL pipeline)。 触发场景:创建数据集成任务 / 数据同步任务 / 管道任务 / pipeline / 数据搬运 / reader-writer 配置 / MySQL→MaxCompute / Doris→PostgreSQL / 离线集成 / create-pipeline / update-pipeline / create-pipeline-node。 覆盖两条路径:两步法(create-pipeline-node 建草稿 → update-pipeline 填配置提交) 和 一步法(create-pipeline)。 关键坑:PluginConfig 必须是 JSON 字符串;columnMappings 顺序敏感且必填;空 PluginConfig 触发 ClassCastException。 触发词:创建管道任务、数据同步、数据集成、数据搬运、pipeline、create-pipeline、update-pipeline、reader writer、PluginConfig、MySQL→MaxCompute、Doris→PostgreSQL。 |
新建集成管道任务 skill
适用场景
- 通过 CLI 创建一个离线集成管道任务(offline pipeline / 实时 / 工作流同理)
- 典型链路:reader (MySQL/Oracle/Doris/...) → writer (MaxCompute/Hive/PostgreSQL/...) 一对一搬运
💡 术语:ODPS(Open Data Processing Service)是 MaxCompute 的旧名称,在 pipeline PluginConfig、API 参数中仍可能出现 odps 字样,均指 MaxCompute。
- 需要把 Steps(reader/writer 插件配置)、Hops(DAG 边)、调度 + 资源 settings 一次性提交
两条 CLI 路径
| 路径 | 命令组合 | 适用 |
|---|
| A. 两步法(推荐) | dev create-pipeline-node(建空草稿) → dev update-pipeline(填配置 + 提交) | 想分阶段:先占名/占目录,再慢慢调试 Steps |
| B. 一步法 | dev create-pipeline(直接带完整 config 创建并提交) | 配置已稳定、CI 化场景 |
共同点:两条路径最终落库的 pipelineDTO.steps[].pluginConfig 结构完全相同;本 skill 的 PluginConfig 参考片段对两者通用。
通用顶层参数
--tenant-id <租户ID> 必填(profile 已配置可省);多租户共享 endpoint 时必须显式传项目所属租户,否则报 DPN.Filter.ProjectNotFound
--project-id <项目ID> 必填(profile 已配置可省)
--env DEV|PROD 仅 update-pipeline / create-pipeline 需要;create-pipeline-node 不需要
路径 A:两步法
A-1. 创建空草稿
aliyun dataphin-public create-pipeline-node \
--tenant-id <tenant-id> \
--project-id <project-id> \
--pipeline-name <task-name> \
--pipeline-type OFFLINE_PIPELINE \
--node-type NORMAL \
--file-info '{"FileName":"<task-name>","Directory":"/"}'
返回:
{
"Data": {
"PipelineId": <int>,
"SubmitId": null,
"Version": null,
"NodeId": null
},
"Code": "OK", "Success": true
}
| 参数 | 说明 |
|---|
--pipeline-type | OFFLINE_PIPELINE / REAL_TIME_PIPELINE |
--node-type | NORMAL / MANUAL / REAL_TIME |
FileInfo.Directory | 默认 /;非 / 必须先存在(否则报错) |
A-2. 填充 Steps/Hops 并提交
aliyun dataphin-public update-pipeline \
--tenant-id <tenant-id> \
--project-id <project-id> \
--env DEV \
--node-info '{"NodeName":"<task-name>","PipelineId":<上一步PipelineId>}' \
--pipeline-config '<见下方 JSON>' \
--schedule-config '<见下方 JSON>' \
--settings '<见下方 JSON>' \
--submit
路径 B:一步法
aliyun dataphin-public create-pipeline \
--tenant-id <tenant-id> \
--project-id <project-id> \
--env DEV \
--pipeline-type 0 \
--mode PIPELINE \
--node-info '{"NodeName":"<task-name>","Directory":"/"}' \
--pipeline-config '<见下方 JSON>' \
--schedule-config '<见下方 JSON>' \
--settings '<见下方 JSON>' \
--submit
--pipeline-type 取值:0 = 离线集成(默认) / 1 = 实时 / 14 = 工作流。
--pipeline-config 完整骨架
pipeline-config 是集成任务最复杂的字段(含 reader/writer/transformer/column 映射),按 reader/writer 类型组合的完整骨架抽离到独立 reference:
📖 详见 references/pipeline-config.md(涵盖 MySQL→MaxCompute / Doris→PG / Oracle→Hive 等常用组合)
关键规则速查:
- CLI 的 autocreate 不生效:目标表必须手动预建
- 类型映射陷阱:Doris LARGEINT→PG NUMERIC、Doris TINYINT→PG SMALLINT
- column 顺序:reader.column 与 writer.column 必须一一对应、长度一致
--schedule-config 完整骨架
{
"ScheduleType": "NORMAL",
"CronExpression": "0 0 0 * * ?",
"ScheduleStartTime": "1970-01-01 00:00:00",
"ScheduleEndTime": "9999-01-01 00:00:00",
"ScheduleIntervalType": "DAILY",
"ReRunMode": "ALL_ALLOWED",
"NodeStatus": 1,
"Priority": 5,
"ResourceGroupId": "default",
"DevResourceGroupId": "default",
"ExecuteTimeOutConfig": { "FollowSystem": true
--settings 完整骨架
{
"RequiredResource": { "Cpus": 0.5, "MemoryInMb": 1024 },
"JvmOption": "",
"NoFlowTimeout": 30,
"Engine": { "Name": "dlink" },
"ErrorLimit": { "Record": 0 },
"TimeZone": "Asia/Shanghai",
"SqlTimeout": 30,
"Speed": { "Concurrent": 3 },
"ConnectRetryTime": [
{
校验
aliyun dataphin-public get-pipeline-by-id \
--tenant-id <tenant-id> \
--project-id <project-id> \
--env DEV \
--pipeline-id <pipelineId>
只建草稿没填 Steps 时,Data 可能为 null(预期);填好 Steps 后再查应返回完整 Steps / Hops / Settings。
常见坑
PluginConfig 必须是 JSON 字符串:CLI 不会递归序列化嵌套对象。把每个插件 config 用 JSON.stringify 转字符串后再放进 Steps[].PluginConfig。
- 空
PluginConfig: "{}" 触发 ClassCastException:服务端反序列化为 DefaultOutputPluginConfig 与 BaseOutputPluginConfig 类型转换失败。最少要带 dsName/dsId/dsType/table/columns。
--tenant-id 必须与项目所属租户一致:多租户共享同一 endpoint 时,profile 中的 tenant_id 与目标项目的租户可能不同,必须显式传项目租户,否则 DPN.Filter.ProjectNotFound。
columnMappings 必填且顺序敏感:MaxCompute writer 必须显式声明每一列的 sourceColumn → targetColumn,inputColumnIndex 从 0 开始且与 reader columns 顺序对齐,否则跑批数据错位。
- 大整数 ID 字符串化:
dsId / nodeId / fileId 体量超 Number.MAX_SAFE_INTEGER(如 7445807200604583744)必须以字符串传入,避免 JS JSON.parse 精度丢失。
- 缺省上游需挂租户虚拟根节点:
UpStreamList 不能为空,否则提交时报 NodeWithoutUpstream。租户虚拟根节点的查找见 find-tenant-root-node(经套件入口路由加载)。
Directory 必须已存在:默认 / 永远存在;自定义目录前需先建好对应类型为 offlinePipeline 的目录。
prodTableNotExistAction: autocreate 在 CLI 不生效:create-pipeline / update-pipeline 走 OpenAPI 路径时会先校验目标表存在(DPN.Os.TableNotFound),即使配置了 autocreate 也不行。必须先用 execute-ad-hoc-task --operator-type MaxCompute_SQL 执行 DDL 建好目标表,再提交 pipeline 配置。
schedule-config 的 cron 字段名是 CronExpression(不是 ScheduleCron),用错会报「调度周期表达式为空」。
- 查询 MySQL 源表字段用
execute-ad-hoc-task --operator-type DATABASE_SQL:MySQL/Oracle/PostgreSQL/SQLServer 等关系型数据库统一使用 ,必须同时传 和 。查询结果在 (从 开始)的 字段中。
完整示例 1(路径 A,MySQL→MaxCompute 一对一搬运)
TASK_ID=$(aliyun dataphin-public execute-ad-hoc-task \
--tenant-id <tenant-id> \
--project-id <project-id> \
--operator-type DATABASE_SQL \
--data-source-id <mysql-ds-id> \
--data-source-schema <db-name> \
--code "SELECT COLUMN_NAME, DATA_TYPE FROM information_schema.columns WHERE table_schema='<db>' AND table_name='<table>' ORDER BY ORDINAL_POSITION" \
--output json | jq -r '.ExecuteResult.TaskId')
aliyun dataphin-public get-ad-hoc-task-result \
--tenant-id <tenant-id> \
--project-id <project-id> \
--task-id "${TASK_ID}" \
--sub-task-id 0
aliyun dataphin-public execute-ad-hoc-task \
--tenant-id <tenant-id> \
--project-id <project-id> \
--operator-type MaxCompute_SQL \
--code "CREATE TABLE IF NOT EXISTS <table> (<col1> string, <col2> double, ...) PARTITIONED BY (ds string) LIFECYCLE 3600;"
PID=$(aliyun dataphin-public create-pipeline-node \
--tenant-id <tenant-id> \
--project-id <project-id> \
--pipeline-name <task-name> \
--pipeline-type OFFLINE_PIPELINE \
--node-type NORMAL \
--file-info '{"FileName":"<task-name>","Directory":"/"}' \
--output json | jq -r '.Data.PipelineId')
aliyun dataphin-public update-pipeline \
--tenant-id <tenant-id> \
--project-id <project-id> \
--env DEV \
--node-info "{\"NodeName\":\"<task-name>\",\"PipelineId\":${PID}}" \
--pipeline-config "$(cat pipeline-config.json)" \
--schedule-config "$(cat schedule-config.json)" \
--settings "$(cat settings.json)" \
--submit
aliyun dataphin-public get-pipeline-by-id \
--tenant-id <tenant-id> \
--project-id <project-id> \
--env DEV \
--pipeline-id "${PID}"
完整示例 2(路径 A,Doris→PostgreSQL 一对一搬运)
aliyun dataphin-public execute-ad-hoc-task \
--tenant-id <tenant-id> \
--project-id <project-id> \
--operator-type DATABASE_SQL \
--data-source-id <pg-ds-id> \
--data-source-schema <schema-name> \
--code 'CREATE TABLE IF NOT EXISTS demo02 (
user_id NUMERIC NOT NULL,
username VARCHAR(50) NOT NULL,
city VARCHAR(20),
age SMALLINT,
sex SMALLINT,
PRIMARY KEY (user_id, username)
);'
PID=$(aliyun dataphin-public create-pipeline-node \
--tenant-id <tenant-id> \
--project-id <project-id> \
--pipeline-name <task-name> \
--pipeline-type OFFLINE_PIPELINE \
--node-type NORMAL \
--file-info '{"FileName":"<task-name>","Directory":"/"}' \
--output json | jq -r '.Data.PipelineId')
aliyun dataphin-public update-pipeline \
--tenant-id <tenant-id> \
--project-id <project-id> \
--env DEV \
--node-info "{\"NodeName\":\"<task-name>\",\"PipelineId\":${PID}}" \
--pipeline-config "$(cat pipeline-config.json)" \
--schedule-config "$(cat schedule-config.json)" \
--settings "$(cat settings.json)" \
--submit
aliyun dataphin-public get-pipeline-by-id \
--tenant-id <tenant-id> \
--project-id <project-id> \
--env DEV \
--pipeline-id "${PID}"
PostgreSQL 目标端踩坑速查:
schemaName 必填——PG 有 schema 概念(常见 public 或自定义),不传会报表不存在
- 目标表必须手动预建——
prodTableNotExistAction: autocreate 在 CLI/OpenAPI 场景不生效
- 类型不能照搬源端——Doris
LARGEINT 在 PG 不存在,需映射为 NUMERIC;TINYINT 需映射为 SMALLINT
columnMappings[].originalType 填 PG 目标类型,不是 Doris 源类型
- PostgreSQL 建表用
--operator-type DATABASE_SQL(关系型数据库统一入口),必须同时传 --data-source-id 和 --data-source-schema