大数据诊断平台Compass
Compass 简介
Compass 是一个用于诊断大数据生态系统中计算引擎和调度器的平台,旨在提高故障排除的效率并降低问题调整的复杂性。它自动收集日志和指标,使用启发式规则识别问题并提供调整建议。对于日志,Compass 还使用 ChatGPT 提供诊断建议。日志通过 drain 算法自动聚合为模板,可用于人工干预等,提升诊断自动化和优化方案的能力。
主要功能特性
- 非侵入式,即时诊断:无需修改已有的调度平台,即可体验诊断效果。
- 支持多种主流调度平台:例如 DolphinScheduler 2.x 和 3.x、Airflow 或自研平台等。
- 支持多版本 Spark、MapReduce、Flink、Hadoop 2.x 和 3.x 任务日志诊断和解析。
- 支持工作流层异常诊断:识别各种失败和基线耗时异常问题。
- 支持引擎层异常诊断:包含数据倾斜、大表扫描、内存浪费等 14 种异常类型。
- 支持各种日志匹配规则编写和异常阈值调整:可根据实际场景优化。
- 支持一键诊断全量(包含非调度平台提交任务)Spark/MapReduce 任务。
- 支持 ChatGPT 对异常日志进行诊断:提供解决方案,使用 drain 算法聚合模板,节约成本。
支持组件
- ChatGPT
- Spark
- Flink
- MapReduce
- Trino
- Spark Tez
- Airflow
- DolphinScheduler
- Azkaban
- Oozie
- Debezium(同步 PostgreSQL 到 PostgreSQL 的数据同步)
- 其他(我们非常欢迎与倾听其他任务建设性意见)
社区
欢迎加入社区咨询使用或成为 Compass 开发者。以下是获得帮助的方法:
- 提交 issue。
- 提交 pull request,请阅读 contributing guideline。
- 讨论 Idea & Question。
我们将会尽快回复。
罗盘已支持诊断类型概览
| 引擎 | 诊断维度 | 诊断类型 | 类型说明 |
|---|---|---|---|
| Spark | 失败分析 | 运行失败 | 最终运行失败的任务 |
| 首次失败 | 重试次数大于 1 的成功任务 | ||
| 长期失败 | 最近 10 天运行失败的任务 | ||
| 耗时分析 | 基线时间异常 | 相对于历史正常结束时间,提前结束或晚点结束的任务 | |
| 基线耗时异常 | 相对于历史正常运行时长,运行时间过长或过短的任务 | ||
| 运行耗时长 | 运行时间超过 2 小时的任务 | ||
| 报错分析 | SQL 失败 | 因 SQL 执行问题而导致失败的任务 | |
| Shuffle 失败 | 因 Shuffle 执行问题而导致失败的任务 | ||
| 内存溢出 | 因内存溢出问题而导致失败的任务 | ||
| 资源分析 | 内存浪费 | 内存使用峰值与总内存占比过低的任务 | |
| CPU 浪费 | driver/executor 计算时间与总 CPU 计算时间占比过低的任务 | ||
| 效率分析 | 大表扫描 | 没有限制分区导致扫描行数过多的任务 | |
| OOM 预警 | 广播表的累计内存与 driver 或 executor 任意一个内存占比过高的任务 | ||
| 数据倾斜 | stage 中存在 task 处理的最大数据量远大于中位数的任务 | ||
| Job 耗时异常 | job 空闲时间与 job 运行时间占比过高的任务 | ||
| Stage 耗时异常 | stage 空闲时间与 stage 运行时间占比过高的任务 | ||
| Task 长尾 | stage 中存在 task 最大运行耗时远大于中位数的任务 | ||
| HDFS 卡顿 | stage 中存在 task 处理速率过慢的任务 | ||
| 推测执行 Task 过多 | stage 中频繁出现 task 推测执行的任务 | ||
| 全局排序异常 | 全局排序导致运行耗时过长的任务 | ||
| MapReduce | 资源分析 | 内存浪费 | 内存使用峰值与总内存占比过低的任务 |
| 效率分析 | 大表扫描 | 扫描行数过多的任务 | |
| Task 长尾 | map/reduce task 最大运行耗时远大于中位数的任务 | ||
| 数据倾斜 | map/reduce task 处理的最大数据量远大于中位数的任务 | ||
| 推测执行 Task 过多 | map/reduce task 中频繁出现推测执行的任务 | ||
| GC 异常 | GC 时间相对 CPU 时间占比过高的任务 | ||
| Flink | 资源诊断 | 内存利用率高 | 计算内存的使用率,如果使用率高于阈值,则增加内存 |
| 内存利用率低 | 计算内存的使用率,如果使用率低于阈值,则降低内存 | ||
| JM 内存优化 | 根据 tm 个数计算 jm 内存的建议值 | ||
| 作业无流量 | 检测作业的 Kafka source 算子是否没有流量 | ||
| TM 管理内存优化 | 计算作业管理内存的使用率,给出合适的管理内存建议值 | ||
| 部分 TM 空跑 | 检测是否有 tm 没有流量,并且 CPU 和内存也没有使用 | ||
| 并行度不够 | 检测作业是否因为并行度不够引起延迟 | ||
| CPU 利用率高 | 计算作业的 CPU 均值使用率,如果高于阈值,则增加 CPU | ||
| CPU 利用率低 | 计算作业的 CPU 均值使用率,如果低于阈值,则降低 CPU | ||
| CPU 峰值利用率高 | 计算作业的 CPU 峰值使用率,如果高于阈值,则增加 CPU | ||
| 异常诊断 | 存在慢算子 | 检测作业是否存在慢算子 | |
| 存在反压算子 | 检测作业是否存在反压算子 | ||
| 作业延迟高 | 检测作业的 Kafka 延迟是否高于阈值 |
UI
Spark







Flink




系统架构
系统架构图


架构说明
整体架构分 3 层:
- 调度系统对接层:实现对接调度器、Yarn、Spark、Flink、HDFS 等系统,同步任务及其日志元数据到诊断系统。
- 服务层:包括数据采集、元数据关联&模型标准化、异常检测、资源诊断、Portal 模块。
- 基础组件层:包括 MySQL、OpenSearch、Kafka、Redis、Zookeeper 等组件。
具体模块流程阶段:
- 数据采集阶段:task-canal/adapter 模块订阅同步调度系统的用户、DAG、作业、执行记录等工作流元数据同步至诊断平台;task-metadata 模块定时同步 Yarn ResourceManager、Spark HistoryServer App 元数据至诊断系统,关联日志存储路径,为后续数据处理阶段作基础。
- 数据关联与模型标准化阶段:task-syncer 模块将同步的数据标准化为 User、Project、Flow、Task、TaskInstance 模型;task-application 模块将工作流层与引擎层元数据关联。
- 工作流层&引擎层异常检测阶段:至此已经获得数据标准模型,针对标准模型进一步 Workflow 异常检测流程。task-detect 模块进行工作流层异常任务检测,例如运行失败、基线耗时异常等;task-parser 模块进行引擎层异常任务检测,例如 SQL 失败、Shuffle 失败等;task-flink 模块进行 Flink 作业资源及异常检测,例如 CPU 利用率低,内存利用率低等。
- 业务视图:task-portal 模块提供用户报告总览、一键诊断、工作流层任务诊断、引擎层作业 Application 诊断、诊断建议和详细报告、白名单等功能。
Compass 部署指南
Compass 依赖了 Canal、PostgreSQL(或 MySQL)、Kafka、Redis、Zookeeper、OpenSearch,需要提前准备好相关环境。
环境要求
| Dependency | Version | Optional | Description |
|---|---|---|---|
| Canal | v1.1.6+ | yes | needed by Airflow, DolphinScheduler |
| MySQL | 5.7+ | yes | |
| PostgreSQL | 10.0+ | no | |
| Kafka | all | no | |
| Redis | all | no | deployed in cluster mode |
| Zookeeper | 3.4.5 | no | needed by canal |
| OpenSearch | 1.3.12 | no |
OpenSearch 兼容 Elasticsearch 7.0+。
Compass 支持单机和集群部署,可按模块弹性扩缩容。
编译
请使用 JDK 8 以及 maven 3.6.0+ 进行编译,构建流程步骤如下:
bash
git clone https://github.com/cubefs/compass.git
cd compass
mvn clean package -DskipTests -Pdist
或者
bash
mvn clean package -DskipTests -Pdist,spark # 打包时 web 展示只有 spark 诊断页面
或者
bash
mvn clean package -DskipTests -Pdist,flink # 打包时 web 展示只有 flink 诊断页面
使用 docker compose 启动应用服务:
bash
cp dist/compass-v1.1.2.tar.gz docker/playground
cd docker/playground/
docker compose --profile dependencies up -d
docker compose --profile compass-demo up -d
更多请查看文档 docker-playground。
初始化数据库
支持 PostgreSQL(默认)或 MySQL(需要手动下载 mysql-connector-java 复制到各模块 lib 目录下,canal 除外)作为元数据存储。表结构由两部分组成:document/sql/compass_*.sql,document/sql/dolphinscheduler_*.sql(需要根据实际使用版本修改,支持 2.x 和 3.x)或 document/sql/airflow_*.sql(支持 2.x)。如果您使用的是自研调度平台,请参考上述 SQL 表结构。
修改配置和启动
compass/bin 和 compass/conf 是作为公共脚本和配置使用,方便统一启停和配置管理。
bash
# 启动所有模块
./bin/start_all.sh
# 停止所有模块
./bin/stop_all.sh
compass_env.sh 配置说明:
Kafka 需要预先创建好 topic: mysqldata, task-instance, task-application, exception-log
bash
#!/bin/bash
# dolphinscheduler or airflow or custom
export SCHEDULER="dolphinscheduler"
export SPRING_PROFILES_ACTIVE="hadoop,${SCHEDULER}"
# Configuration for Scheduler MySQL, compass will subscribe data from scheduler database via canal
export SCHEDULER_MYSQL_ADDRESS="localhost:3306"
export SCHEDULER_MYSQL_DB="dolphinscheduler"
export SCHEDULER_DATASOURCE_URL="jdbc:mysql://${SCHEDULER_MYSQL_ADDRESS}/${SCHEDULER_MYSQL_DB}?useUnicode=true&characterEncoding=utf-8&serverTimezone=Asia/Shanghai"
export SCHEDULER_DATASOURCE_USERNAME=""
export SCHEDULER_DATASOURCE_PASSWORD=""
# Configuration for compass database (mysql or postgresql)
export DATASOURCE_TYPE="mysql"
export COMPASS_DATASOURCE_ADDRESS="localhost:3306"
export COMPASS_DATASOURCE_DB="compass"
export SPRING_DATASOURCE_URL="jdbc:${DATASOURCE_TYPE}://${COMPASS_DATASOURCE_ADDRESS}/${COMPASS_DATASOURCE_DB}"
export SPRING_DATASOURCE_USERNAME=""
export SPRING_DATASOURCE_PASSWORD=""
# Configuration for compass Kafka, used to subscribe data by canal and log queue, etc. (default version: 3.4.0)
export SPRING_KAFKA_BOOTSTRAPSERVERS="host1:port,host2:port"
# Configuration for compass redis, used to cache and log queue, etc. (cluster mode)
export SPRING_REDIS_CLUSTER_NODES="localhost:6379"
# Optional
export SPRING_REDIS_PASSWORD=""
# Zookeeper (cluster: 3.4.5, needed by canal)
export SPRING_ZOOKEEPER_NODES="localhost:2181"
# OpenSearch (default version: 1.3.12) or Elasticsearch (7.x~)
export SPRING_OPENSEARCH_NODES="localhost:9200"
# Optional
export SPRING_OPENSEARCH_USERNAME=""
# Optional
export SPRING_OPENSEARCH_PASSWORD=""
# Optional, needed by OpenSearch, keep empty if OpenSearch does not use truststore.
export SPRING_OPENSEARCH_TRUSTSTORE=""
# Optional, needed by OpenSearch, keep empty if OpenSearch does not use truststore.
export SPRING_OPENSEARCH_TRUSTSTOREPASSWORD=""
# Prometheus for flink, ignore it if you do not need flink.
export FLINK_PROMETHEUS_HOST="http://localhost:9090"
export FLINK_PROMETHEUS_TOKEN=""
export FLINK_PROMETHEUS_DATABASE=""
# Optional, needed by task-gpt module to get exception solution, ignore if you do not need it.
export CHATGPT_ENABLE=false
# Openai keys needed by enabling chatgpt, random access the key if there are multiple keys.
export CHATGPT_API_KEYS=sk-xxx1,sk-xxx2
# Optional, needed if setting proxy, or keep it empty.
export CHATGPT_PROXY="" # for example, https://proxy.ai
# chatgpt model
export CHATGPT_MODEL="gpt-3.5-turbo"
# chatgpt prompt
export CHATGPT_PROMPT="You are a senior expert in big data, teaching beginners. I will give you some anomalies and you will provide solutions to them."
# task-canal 模块配置
# 调度平台 MySQL 订阅账号,确定是否已开启 binlog
export CANAL_INSTANCE_MASTER_ADDRESS=${SCHEDULER_MYSQL_ADDRESS}
export CANAL_INSTANCE_DBUSERNAME=${SCHEDULER_DATASOURCE_USERNAME}
export CANAL_INSTANCE_DBPASSWORD=${SCHEDULER_DATASOURCE_PASSWORD}
# 需要订阅的库表配置过滤
if [ ${SCHEDULER} == "dolphinscheduler" ]; then
export CANAL_INSTANCE_FILTER_REGEX="${SCHEDULER_MYSQL_DB}.t_ds_user,${SCHEDULER_MYSQL_DB}.t_ds_project,${SCHEDULER_MYSQL_DB}.t_ds_task_definition,${SCHEDULER_MYSQL_DB}.t_ds_task_instance,${SCHEDULER_MYSQL_DB}.t_ds_process_definition,${SCHEDULER_MYSQL_DB}.t_ds_process_instance,${SCHEDULER_MYSQL_DB}.t_ds_process_task_relation"
elif [ ${SCHEDULER} == "airflow" ]; then
export CANAL_INSTANCE_FILTER_REGEX="${SCHEDULER_MYSQL_DB}.dag,${SCHEDULER_MYSQL_DB}.serialized_dag,${SCHEDULER_MYSQL_DB}.ab_user,${SCHEDULER_MYSQL_DB}.dag_run,${SCHEDULER_MYSQL_DB}.task_instance"
else
export CANAL_INSTANCE_FILTER_REGEX=".*\..*"
fi
application-hadoop.yml 说明:
yaml
hadoop:
# task-applicaiton & task-parser 模块配置依赖
namenodes:
- nameservices: logs-hdfs # dfs.nameservices 属性值
namenodesAddr: [ "machine1.example.com", "machine2.example.com" ] # dfs.namenode.rpc-address.[nameservice ID].[name node ID] 属性值
namenodes: ["nn1", "nn2"] # dfs.ha.namenodes.[nameservice ID] 属性值
user: hdfs # 用户
password: # 密码,如果没开启鉴权,则不需要
port: 8020 # 端口
matchPathKeys: [ "flume" ] # task-application 模块使用,调度平台日志 hdfs 路径关键字
# kerberos
enableKerberos: false
# /etc/krb5.conf
krb5Conf: ""
# hdfs/*@EXAMPLE.COM
principalPattern: ""
# admin
loginUser: ""
# /var/kerberos/krb5kdc/admin.keytab
keytabPath: ""
# task-metadata 模块配置依赖
yarn:
- clusterName: "bigdata"
resourceManager: [ "ip:port" ] # yarn.resourcemanager.webapp.address 属性值
jobHistoryServer: "ip:port" # mapreduce.jobhistory.webapp.address 属性值
spark:
sparkHistoryServer: [ "ip:port" ] # spark history ui 地址
Compass 模块介绍
工程目录
plaintext
compass
├── bin
│ ├── compass_env.sh 环境变量,基础组件配置
│ ├── start_all.sh 启动脚本
│ └── stop_all.sh 停止脚本
├── conf
│ └── application-hadoop.yml hadoop 相关配置
├── task-application 关联任务实例、applicationId、hdfs_log_path
├── task-canal 订阅调度平台 MySQL 表元数据到 Kafka
├── task-canal-adapter 同步调度平台 MySQL 表元数据 Compass 平台
├── task-detect 工作流层异常类型检测
├── task-metadata 同步 Yarn、Spark 任务元数据到 OpenSearch
├── task-parser 日志解析和 Spark 任务异常检测
├── task-portal 异常任务的可视化服务
├── task-flink Flink 任务资源及异常诊断
├── task-flink-core Flink 任务诊断规则逻辑
├── task-portal 异常任务的可视化服务
├── task-gpt 聚合日志模板,并使用 chatgpt 给模板解决方案
└── task-syncer 调度平台任务关系表的抽象和映射
历史数据同步
如果调度平台的数据库是 MySQL,Compass 数据库是 PostgreSQL,可使用 pgloader 创建依赖表和同步历史全量数据。
同步 dolphinscheduler 表:
sql
LOAD DATABASE
FROM mysql://root:password@localhost:3306/dolphinscheduler
INTO postgresql://postgres@localhost:5432/compass
ALTER SCHEMA 'dolphinscheduler' RENAME TO 'public'
INCLUDING ONLY TABLE NAMES MATCHING 't_ds_process_definition','t_ds_process_instance','t_ds_process_task_relation','t_ds_project','t_ds_task_definition','t_ds_task_instance','t_ds_user';
同步 airflow 表:
sql
LOAD DATABASE
FROM mysql://root:password@localhost:3306/airflow
INTO postgresql://postgres@localhost:5432/compass_airflow
ALTER SCHEMA 'airflow' RENAME TO 'public'
INCLUDING ONLY TABLE NAMES MATCHING 'task_instance','dag_run','ab_user','dag','serialized_dag'
ALTER TABLE NAMES MATCHING 'task_instance' RENAME TO 'tb_task_instance'
ALTER TABLE NAMES MATCHING 'dag_run' RENAME TO 'tb_dag_run'
ALTER TABLE NAMES MATCHING 'ab_user' RENAME TO 'tb_ab_user'
ALTER TABLE NAMES MATCHING 'dag' RENAME TO 'tb_dag'
ALTER TABLE NAMES MATCHING 'serialized_dag' RENAME TO 'tb_serialized_dag';
如果调度平台的数据库是 MySQL,Compass 数据库是 MySQL,可使用 task-canal-adapter 接口进行同步历史全量数据。
同步 dolphinscheduler 表:
bash
curl "localhost:8181/etl/rdb/mysql1/t_ds_process_definition.yml" -X POST
curl "localhost:8181/etl/rdb/mysql1/t_ds_process_instance.yml" -X POST
curl "localhost:8181/etl/rdb/mysql1/t_ds_process_task_relation.yml" -X POST
curl "localhost:8181/etl/rdb/mysql1/t_ds_project.yml" -X POST
curl "localhost:8181/etl/rdb/mysql1/t_ds_task_definition.yml" -X POST
curl "localhost:8181/etl/rdb/mysql1/t_ds_task_instance.yml" -X POST
curl "localhost:8181/etl/rdb/mysql1/t_ds_user.yml" -X POST
同步 airflow 表:
bash
curl "localhost:8181/etl/rdb/mysql1/airflow_db_ab_user.yml" -X POST
curl "localhost:8181/etl/rdb/mysql1/airflow_db_dag_run.yml" -X POST
curl "localhost:8181/etl/rdb/mysql1/airflow_db_dag.yml" -X POST
curl "localhost:8181/etl/rdb/mysql1/airflow_db_task_instance.yml" -X POST
task-canal
如果您使用的是 DolphinScheduler 或者 Airflow 或者自研的等调度平台,元数据存储在 MySQL,可使用 canal.deployer 订阅 MySQL binlog 同步到 Kafka,默认 topic 是 mysqldata。
plaintext
task-canal
├── bin
│ ├── compass_env.sh compass 环境变量
│ ├── init_canal.sh 下载 canal.deployer 依赖包
│ ├── restart.sh
│ ├── startup.sh
│ └── stop.sh
├── canal.deployer-1.1.6.tar.gz compass 不提供 canal 依赖包,可通过 init_canal.sh 下载,若无网络则自行下载到 task-canal 根目录
├── conf
│ ├── example
│ │ ├── instance.properties 源 MySQL 配置和库表配置
│ ├── canal_local.properties zk, kafka 等配置
│ ├── canal.properties
│ ├── logback.xml
├── lib
└── plugin
核心配置
conf/example/instance.properties
properties
canal.instance.master.address=localhost:33066
canal.instance.dbUsername=root
canal.instance.dbPassword=root
canal.instance.filter.regex=.*\..*
canal.mq.topic=mysqldata
# 动态 topic 和分区默认不配置,若数据量比较大,可按表 Hash 到相同 topic 不同分区,避免单分区压力过大
canal.mq.dynamicTopic = mysqldata:db\..*
canal.mq.partitionsNum = 12
canal.mq.partitionHash = .*\..*
conf/canal.properties
properties
canal.zkServers = localhost:2181
canal.serverMode = kafka
kafka.bootstrap.servers = localhost:9092
task-canal-adapter
canal.adapter 模块作用:同步依赖调度平台的元数据表到 compass,只同步任务相关和用户表,其他表按需同步。例如对于 DolphinScheduler:t_ds_project.yml 定义同步了 ds_project 表,若需要同步其他表可参考 conf/rdb 下配置。项目中已提供 DolphinScheduler 和 Airflow 平台同步模板,若使用其他平台可参考模板。
plaintext
task-canal-adapter
├── bin
│ ├── compass_env.sh
│ ├── init_canal_adapter.sh 下载 canal.adapter 压缩包和解压相关 lib 和 plugin
│ ├── restart.sh
│ ├── startup.sh
│ └── stop.sh
├── canal.adapter-1.1.6.tar.gz compass 不提供 canal 依赖包,可通过 init_canal.sh 下载,若无网络则自行下载到 task-canal-adapter 根目录
├── conf
│ ├── application.yml
│ └── rdb
│ ├── airflow_db_ab_user.yml
│ ├── airflow_db_dag_run.yml
│ ├── airflow_db_dag.yml
│ ├── airflow_db_task_instance.yml
│ ├── t_ds_process_definition.yml
│ ├── t_ds_process_instance.yml
│ ├── t_ds_process_task_relation.yml
│ ├── t_ds_project.yml
│ ├── t_ds_task_definition.yml
│ ├── t_ds_task_instance.yml
│ ├── t_ds_user.yml
│ └── template.yml
├── lib
└── plugin
表数据全量同步接口
示例:curl "localhost:8181/etl/rdb/mysql1/template.yml" -X POST
其中 template.yml 即为 conf/rdb 下的配置文件。
核心配置
conf/application.yml
yaml
canal.conf:
srcDataSources:
defaultDS:
# Scheduling platform MySQL synchronization account
url: ${CANAL_ADAPTER_SOURCE_DATASOURCE_URL}
username: ${CANAL_ADAPTER_SOURCE_DATASOURCE_USERNAME}
password: ${CANAL_ADAPTER_SOURCE_DATASOURCE_PASSWORD}
canalAdapters:
- instance: mysqldata # kafka topic
groups:
- groupId: g1
outerAdapters:
- name: rdb
key: mysql1
properties:
# Compass platform datasource account
jdbc.url: ${CANAL_ADAPTER_DESTINATION_DATASOURCE_URL}
jdbc.username: ${CANAL_ADAPTER_DESTINATION_DATASOURCE_USERNAME}
jdbc.password: ${CANAL_ADAPTER_DESTINATION_DATASOURCE_PASSWORD}
conf/rdb/template.yml
yaml
dataSourceKey: defaultDS
destination: mysqldata
groupId: g1
outerAdapterKey: mysql1
concurrent: false
dbMapping:
database: ${SCHEDULER_MYSQL_DB}
# 调度平台 MySQL 表
table: example
# compass 平台 MySQL 表
targetTable: example
# 主键配置
targetPk:
id: id
mapAll: true
commitBatch: 1
task-syncer
task-syncer 模块是关联调度平台和 compass 的抽象层,使得 compass 能够兼容和诊断不同的调度平台任务。该模块抽象定义了 compass 核心依赖的关系表:
user:登录和权限校验,隔离不同用户权限project:项目关系flow:工作流定义关系task:具体任务定义关系task_instance:任务运行实例
其中关系是 user -> project -> flow -> task -> task_instance,可根据实际调度平台自行定义关系。
plaintext
task-syncer
├── bin
│ ├── compass_env.sh
│ ├── startup.sh
│ └── stop.sh
├── conf
│ ├── application-airflow.yml
│ ├── application-dolphinscheduler.yml
│ ├── application.yml
│ └── logback.xml
├── lib
核心配置
conf/application-xxx.yml 定义了数据同步表字段之间的映射关系,实现源表和目标表的转化。
columnMapping实现了字段之间的映射columnValueMapping实现了字段值的映射constantColumn实现了常量列的映射columnDep实现了列字段值依赖查询,可自定义 SQL 实现表字段之间的关联
下面以同步 DolphinScheduler 调度平台示例说明。
user 表映射:
yaml
# DolphinScheduler 库名
- schema: "dolphinscheduler"
# DolphinScheduler user 表
table: "t_ds_user"
# compass user 表
targetTable: "user"
# columnMapping 用于字段映射,key 是 compass 定义的字段,值是 DolphinScheduler 定义的字段
columnMapping:
user_id: "id"
username: "user_name"
password: "user_password"
is_admin: "user_type"
email: "email"
phone: "phone"
create_time: "create_time"
update_time: "update_time"
# 字段值映射, 目标字段值, 源字段值, 字段类型
columnValueMapping:
is_admin: [ { targetValue: "0", originValue: [ "0" ] }, { targetValue: "1", originValue: [ "1" ] } ]
# 常量列定义
constantColumn:
scheduler_type: "DolphinScheduler"
task_instance 表映射:
yaml
- schema: "dolphinscheduler"
table: "t_ds_task_instance"
targetTable: "task_instance"
columnMapping:
id: "id"
project_name: ""
flow_name: ""
task_name: "name"
start_time: "start_time"
end_time: "end_time"
execution_time: ""
task_state: "state"
task_type: "task_type"
retry_times: "retry_times"
max_retry_times: "max_retry_times"
worker_group: "worker_group"
create_time: "create_time"
update_time: "update_time"
columnValueMapping:
task_state:
- { targetValue: "success", originValue: [ "7", "14" ] }
- { targetValue: "fail", originValue: [ "6", "9" ] }
- { targetValue: "other", originValue: [ "0", "1", "2", "3", "4", "5", "8", "10", "11", "12", "13" ] }
columnDep:
# 列字段值依赖,由于该表缺失了 project_name, flow_name, execution_time 字段,因此需要关联其他表查询
columns: [ "project_name", "flow_name", "execution_time" ]
queries: [ "select t2.schedule_time as execution_time, t3.name as flow_name, t4.name as project_name from t_ds_task_instance as t1 inner join t_ds_process_instance as t2 on t1.process_instance_id = t2.id inner join t_ds_process_definition as t3 on t2.process_definition_code = t3.code inner join t_ds_project as t4 on t3.project_code=t4.code where t1.id=${id}" ]
task-application
task-application 模块关联 task_name、applicationId、hdfs_log_path,该模块需要读取调度平台日志,推荐使用 flume 收集到 hdfs,方便统一做日志诊断和分析。
plaintext
task-application/
├── bin
│ ├── compass_env.sh
│ ├── startup.sh
│ └── stop.sh
├── conf
│ ├── application-airflow.yml
│ ├── application-dolphinscheduler.yml
│ ├── application-hadoop.yml
│ ├── application.yml
│ └── logback.xml
├── lib
核心配置
conf/application-hadoop.yml
yaml
hadoop:
namenodes:
- nameservices: logs-hdfs
namenodesAddr: [ "host1", "host2" ]
namenodes: ["namenode1", "namenode2"]
user: hdfs
password:
port: 8020
matchPathKeys: [ "flume" ]
conf/application-dolphinscheduler/airflow/custom.yml
该配置涉及日志路径规则的拼接,即日志绝对路径的确定。以 flume 收集 dolphinscheduler 到 hdfs 为例,airflow 等同理。表 t_ds_task_instance 记录了日志路径 log_path,但这个是 worker 主机中的目录,上传到 hdfs 的目录
end
