改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

大数据诊断平台Compass

ckckck2025年1月10日49 浏览

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

Spark UI 1
Spark UI 2
Spark UI 3
Spark UI 4
Spark UI 5
Spark UI 6
Spark UI 7

Flink UI 1
Flink UI 2
Flink UI 3
Flink UI 4

系统架构

系统架构图

系统架构图 1
系统架构图 2

架构说明

整体架构分 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_*.sqldocument/sql/dolphinscheduler_*.sql(需要根据实际使用版本修改,支持 2.x 和 3.x)或 document/sql/airflow_*.sql(支持 2.x)。如果您使用的是自研调度平台,请参考上述 SQL 表结构。

修改配置和启动

compass/bincompass/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_nameapplicationIdhdfs_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