Seatunnel 整合 xxl-job 实现批处理任务定时执行
为解决 Seatunnel Web 没有定时任务的问题,本文采用 xxl-job 调度器结合 Seatunnel 实现 批处理的定时执行。
环境要求:在执行器的环境中存在 Seatunnel 环境,并且配置了 Seatunnel 环境变量。
说明:任务采用 Seatunnel 的./bin/seatunnel.sh -c config.template
的方式实现任务的提交 ,通过 -i 参数实现 参数动态传递
创建任务
编写脚本
将处理任务添加到 Shell 脚本中,具体脚本如下。保存后运行即可。通过修改 config_content
值实现任务的修改
批处理
#!/bin/bashSEATUNNEL_CMD="$SEATUNNEL_HOME/bin/seatunnel.sh"# 定义任务停止时执行的清理操作
exit_func() {# 在这里放入你希望在任务停止时执行的操作,比如释放资源、记录日志等$SEATUNNEL_CMD -can "$JOB_ID"exit;
}# 捕获 SIGINT (Ctrl+C) 和 SIGTERM (手动终止) 信号
trap exit_func SIGINT SIGTERM SIGHUP SIGQUIT SIGKILL# 将配置内容写入变量
config_content=$(cat <<EOL
env {# You can set SeaTunnel environment configuration hereparallelism = 2job.mode = "BATCH"checkpoint.interval = 10000
}source {# This is a example source plugin **only for test and demonstrate the feature source plugin**FakeSource {parallelism = 2result_table_name = "fake"row.num = 16schema = {fields {name = "string"age = "int"}}}# If you would like to get more information about how to configure SeaTunnel and see full list of source plugins,# please go to https://seatunnel.apache.org/docs/connector-v2/source
}sink {Console {}
}
EOL
)echo "开始执行任务"## 同步任务,日志中会打印运行日志
echo "$config_content" | $SEATUNNEL_CMD --config /dev/stdin## 异步任务,日志不会记录运行日志# 将配置内容写入标准输入并传递给 SeaTunnel
# SUBJOB_OUTPUT=$(echo "$config_content" | $SEATUNNEL_CMD --config /dev/stdin --async 2>&1)
# #JOB_ID=$(echo "$SUBJOB_OUTPUT" | grep "job name" | awk -F'job id: ' '{print $2}' | awk -F',' '{print $1}')# #echo "任务Id: $JOB_ID"# # 监控任务状态
# while true; do
# # 查询任务状态
# STATUS_OUTPUT=$($SEATUNNEL_CMD -j "$JOB_ID" 2>&1)
# TASK_STATE=$(echo "$STATUS_OUTPUT" | grep "$JOB_ID" | awk -F'"jobStatus":"' '{print $2}' | awk -F '","' '{print $1}')# if [[ "$TASK_STATE" == "FINISHED" ]]; then
# echo "任务完成, 状态: $TASK_STATE"
# exit 0
# fi
# # 检查任务是否已完成
# if [[ "$TASK_STATE" != "RUNNING" ]]; then
# echo "任务已结束,状态:$TASK_STATE"
# exit 1
# else
# echo "任务运行中 ... 状态: $TASK_STATE"
# # 等待 5 秒后再次查询
# sleep 5
# fi
# done
流处理
针对流操作还有部分问题
- 任务无法自启动
- 两次任务无法实现增量同步,每次同步对于Seatunnel来说,都是新任务。
#!/bin/bash
SEATUNNEL_CMD="$SEATUNNEL_HOME/bin/seatunnel.sh"
SEATUNNEL_HOST=localhost
SEATUNNEL_PORT=5801# 定义任务停止时执行的清理操作
exit_func() {# 在这里放入你希望在任务停止时执行的操作,比如释放资源、记录日志等$SEATUNNEL_CMD -can "$JOB_ID"exit;
}# 捕获 SIGINT (Ctrl+C) 和 SIGTERM (手动终止) 信号
trap exit_func SIGINT SIGTERM SIGHUP SIGQUIT SIGKILL# 将配置内容写入变量
config_content=$(cat <<EOL
env {parallelism = 2job.mode = "STREAMING"checkpoint.interval = 2000
}source {FakeSource {parallelism = 2plugin_output = "fake"row.num = 16schema = {fields {name = "string"age = "int"}}}
}sink {Console {}
}
EOL
)echo "开始执行任务"
echo "-------- 配置信息 --------------"
echo "$config_content"
echo "-------- end --------------"# 将配置内容写入标准输入并传递给 SeaTunnel
SUBJOB_OUTPUT=$(echo "$config_content" | $SEATUNNEL_CMD --config /dev/stdin --async 2>&1)
JOB_ID=$(echo "$SUBJOB_OUTPUT" | grep "job name" | awk -F'job id: ' '{print $2}' | awk -F',' '{print $1}')echo "任务Id: $JOB_ID"# 监控任务状态
while true; doSTATUS_OUTPUT=$(curl -s http://$SEATUNNEL_HOST:$SEATUNNEL_PORT/hazelcast/rest/maps/job-info/$JOB_ID)echo $(date "+%Y-%m-%d %H:%M:%S.%3N") "写入数量 : "$(echo "$STATUS_OUTPUT" | awk -F'"SinkWriteCount":"' '{print $2}' | awk -F '","' '{print $1}')", 读取数量 :"$(echo "$STATUS_OUTPUT" | awk -F'"SourceReceivedCount":"' '{print $2}' | awk -F '","' '{print $1}')TASK_STATE=$(echo "$STATUS_OUTPUT" | awk -F'"jobStatus":"' '{print $2}' | awk -F '","' '{print $1}')if [[ "$TASK_STATE" == "FINISHED" ]]; thenecho "任务完成, 状态: $TASK_STATE"exit 0fiif [[ "$TASK_STATE" != "RUNNING" ]]; thenecho "任务已结束,状态:$TASK_STATE"exit 1elseecho "任务运行中 ... 状态: $TASK_STATE"sleep 300fi
done