Flink SQL ジョブを systemd で自動起動させた
はじめに
これまで以下のような構成で、Flink + Kafka + Iceberg を使ったリアルタイム集計基盤を構築してきました。
- syslog / authlog を Kafka に投入
- Flink SQL で Kafka → Iceberg に継続書き込み
- Grafana / Trino / Zeppelin から可視化・分析
- 週次で orphan files cleanup
ただ、実運用してみると次の課題がありました。
- サーバ再起動後に Flink ジョブが自動で復帰しない
この記事では、これらを解決するために、
- Flink クラスタ自動起動
- Flink SQL ジョブ自動投入
- 同名ジョブが既に RUNNING なら submit しない
- systemd による安全な起動管理
- oneshot service の扱い方
までを、実運用向けにまとめます。
前提
想定している環境は以下です。
- Flink standalone cluster
- Kafka topic
syslogauthlog
- Iceberg catalog (Hive Metastore 経由)
- SQL Client で継続ジョブを投入している構成
Flink のインストール先は以下とします。
/opt/flink/current
ジョブ用 SQL は以下に配置します。
/opt/flink/jobs
やりたいこと
最終的にやりたいのはこれです。
- サーバ再起動時に
flink.serviceで Flink クラスタ起動 - その後 systemd が SQL ジョブを自動投入
- ただし 既に同名ジョブが RUNNING / RESTARTING なら何もしない
- つまり ジョブが増殖しない
よくある失敗
最初に結論だけ書くと、以下のような service は危険です。
[Service]
Type=simple
ExecStart=/usr/local/bin/flink-run-sql-job.sh ...
Restart=always
これをやると、
-
sql-client.sh embedded -f ...がジョブを submit - submit 後にプロセスが終了
- systemd が「service が終わった」と判断
-
Restart=alwaysにより再起動 - また submit
となり、同じジョブがどんどん増殖します。
今回の正解構成
今回の正解は次の2点です。
1. systemd 側は oneshot にする
ジョブ service は「起動時に 1 回だけ submit」にする。
2. スクリプト側でも重複起動を防ぐ
flink list -r を見て、同名ジョブが既にあれば submit しない。
この二重ガードでかなり安定します。
ディレクトリ作成
まずは SQL 配置先を作ります。
sudo mkdir -p /opt/flink/jobs
sudo chown -R flink:flink /opt/flink/jobs
SQL ファイルを固定配置する
/tmp/*.sql に置くと、再起動運用で不安定なので、固定パスに置きます。
syslog 用 SQL
sudo tee /opt/flink/jobs/flink_syslog_events.sql >/dev/null <<'EOF'
SET 'table.local-time-zone' = 'Asia/Tokyo';
CREATE CATALOG hive_prod WITH (
'type'='iceberg',
'catalog-type'='hive',
'uri'='thrift://hive1:9083,thrift://hive2:9083',
'warehouse'='hdfs://cluster1/warehouse/iceberg'
);
USE CATALOG hive_prod;
USE logs;
CREATE TEMPORARY TABLE syslog_json (
ts STRING,
src_host STRING,
program STRING,
message STRING
) WITH (
'connector' = 'kafka',
'topic' = 'syslog',
'properties.bootstrap.servers' = 'kafka1:9092,kafka2:9092,kafka3:9092',
'properties.group.id' = 'flink-syslog-events',
'scan.startup.mode' = 'latest-offset',
'format' = 'json',
'json.ignore-parse-errors' = 'true'
);
INSERT INTO hive_prod.logs.syslog_events /*+ OPTIONS(
'compression-codec'='zstd'
) */
SELECT
ts AS ts_raw,
CAST(
TO_TIMESTAMP(
SUBSTRING(ts, 1, 19),
'yyyy-MM-dd''T''HH:mm:ss'
) AS TIMESTAMP(6)
) AS ts,
TRIM(src_host) AS host,
program,
message AS msg,
CAST(SUBSTRING(ts, 1, 10) AS DATE) AS dt,
CAST(SUBSTRING(ts, 12, 2) AS INT) AS hr
FROM syslog_json
WHERE src_host IS NOT NULL
AND TRIM(src_host) <> ''
AND ts IS NOT NULL;
EOF
authlog 用 SQL
sudo tee /opt/flink/jobs/flink_authlog_events.sql >/dev/null <<'EOF'
SET 'table.local-time-zone' = 'Asia/Tokyo';
CREATE CATALOG hive_prod WITH (
'type'='iceberg',
'catalog-type'='hive',
'uri'='thrift://hive1:9083,thrift://hive2:9083',
'warehouse'='hdfs://cluster1/warehouse/iceberg'
);
USE CATALOG hive_prod;
USE logs;
CREATE TEMPORARY TABLE authlog_json (
ts STRING,
src_host STRING,
program STRING,
message STRING
) WITH (
'connector' = 'kafka',
'topic' = 'authlog',
'properties.bootstrap.servers' = 'kafka1:9092,kafka2:9092,kafka3:9092',
'properties.group.id' = 'flink-authlog-events',
'scan.startup.mode' = 'latest-offset',
'format' = 'json',
'json.ignore-parse-errors' = 'true'
);
INSERT INTO hive_prod.logs.authlog_events /*+ OPTIONS(
'compression-codec'='zstd'
) */
SELECT
ts AS ts_raw,
CAST(
TO_TIMESTAMP(
SUBSTRING(ts, 1, 19),
'yyyy-MM-dd''T''HH:mm:ss'
) AS TIMESTAMP(6)
) AS ts,
TRIM(src_host) AS host,
program,
message AS msg,
CAST(SUBSTRING(ts, 1, 10) AS DATE) AS dt,
CAST(SUBSTRING(ts, 12, 2) AS INT) AS hr
FROM authlog_json
WHERE src_host IS NOT NULL
AND TRIM(src_host) <> ''
AND ts IS NOT NULL;
EOF
権限も整えておきます。
sudo chown flink:flink /opt/flink/jobs/flink_syslog_events.sql /opt/flink/jobs/flink_authlog_events.sql
sudo chmod 644 /opt/flink/jobs/flink_syslog_events.sql /opt/flink/jobs/flink_authlog_events.sql
ジョブ投入ラッパースクリプトを作る
ここが今回の肝です。
このスクリプトでは以下をやります。
- Flink CLI 応答待ち
- 既に同名ジョブが動いていれば 何もしない
- 動いていなければ SQL Client で submit
- submit 後に存在確認
/usr/local/bin/flink-run-sql-job.sh
sudo tee /usr/local/bin/flink-run-sql-job.sh >/dev/null <<'EOF'
#!/usr/bin/env bash
set -euo pipefail
JOB_NAME="${1:?job name required}"
SQL_FILE="${2:?sql file required}"
cd /
FLINK_HOME=/opt/flink/current
LOG_DIR=/var/log/flink
LOG_FILE="${LOG_DIR}/${JOB_NAME}.log"
mkdir -p "${LOG_DIR}"
touch "${LOG_FILE}"
log() {
echo "[$(date '+%F %T')] $*" | tee -a "${LOG_FILE}"
}
# Flink CLI が応答するまで待つ
for i in {1..60}; do
if "${FLINK_HOME}/bin/flink" list -r >/dev/null 2>&1; then
break
fi
sleep 2
done
if ! "${FLINK_HOME}/bin/flink" list -r >/dev/null 2>&1; then
log "[ERROR] flink list -r failed"
exit 1
fi
# Running/Restarting に同名ジョブがあれば submit しない
if "${FLINK_HOME}/bin/flink" list -r 2>/dev/null | grep -F " : ${JOB_NAME} (" >/dev/null 2>&1; then
log "[INFO] ${JOB_NAME} already running. skip submit."
exit 0
fi
log "[INFO] submit ${JOB_NAME} using ${SQL_FILE}"
if ! "${FLINK_HOME}/bin/sql-client.sh" embedded -f "${SQL_FILE}" >> "${LOG_FILE}" 2>&1; then
log "[ERROR] sql-client submit failed"
exit 1
fi
sleep 5
# submit 後確認
if "${FLINK_HOME}/bin/flink" list -r 2>/dev/null | grep -F " : ${JOB_NAME} (" >/dev/null 2>&1; then
log "[INFO] ${JOB_NAME} submitted successfully"
exit 0
else
log "[ERROR] ${JOB_NAME} not found after submit"
exit 1
fi
EOF
sudo chmod 755 /usr/local/bin/flink-run-sql-job.sh
sudo chown root:root /usr/local/bin/flink-run-sql-job.sh
systemd service を作る
ジョブ投入は 常駐 service ではなく oneshot にします。
syslog ジョブ用 service
sudo tee /etc/systemd/system/flink-job-syslog-events.service >/dev/null <<'EOF'
[Unit]
Description=Submit Flink SQL Job - syslog_events
After=network-online.target flink.service
Wants=network-online.target
Requires=flink.service
[Service]
Type=oneshot
User=flink
Group=flink
WorkingDirectory=/
Environment=JAVA_HOME=/usr/lib/jvm/java-17-openjdk
Environment=FLINK_CONF_DIR=/etc/flink/conf
Environment=HADOOP_CONF_DIR=/etc/hadoop/conf
Environment=HIVE_CONF_DIR=/etc/hive/conf
ExecStart=/usr/local/bin/flink-run-sql-job.sh insert-into_hive_prod.logs.syslog_events /opt/flink/jobs/flink_syslog_events.sql
RemainAfterExit=yes
TimeoutStartSec=300
[Install]
WantedBy=multi-user.target
EOF
authlog ジョブ用 service
sudo tee /etc/systemd/system/flink-job-authlog-events.service >/dev/null <<'EOF'
[Unit]
Description=Submit Flink SQL Job - authlog_events
After=network-online.target flink.service
Wants=network-online.target
Requires=flink.service
[Service]
Type=oneshot
User=flink
Group=flink
WorkingDirectory=/
Environment=JAVA_HOME=/usr/lib/jvm/java-17-openjdk
Environment=FLINK_CONF_DIR=/etc/flink/conf
Environment=HADOOP_CONF_DIR=/etc/hadoop/conf
Environment=HIVE_CONF_DIR=/etc/hive/conf
ExecStart=/usr/local/bin/flink-run-sql-job.sh insert-into_hive_prod.logs.authlog_events /opt/flink/jobs/flink_authlog_events.sql
RemainAfterExit=yes
TimeoutStartSec=300
[Install]
WantedBy=multi-user.target
EOF
補足: この systemd は「ジョブ監視」ではなく「自動投入」用
今回の systemd サービスは Type=oneshot + RemainAfterExit=yes としているため、
systemd が見ているのは「Flink ジョブの submit が成功したか」までです。そのため、
systemctl status flink-job-xxxx.service
で active (exited) になっていても、それは
submit 自体は成功した
systemd 上は正常終了扱いという意味であり、Flink ジョブ本体が現在も正常稼働していることを保証するものではありません。
実際のジョブ状態確認は Flink 側で行う
以下で確認するのが確実です。
/opt/flink/current/bin/flink list
または Flink Web UI:
http://<flink-host>:8081つまり今回の構成は、
サーバ再起動時に自動でジョブ投入する
人手での再実行を減らすための仕組みであり、
ジョブ死活監視そのものを systemd で行う構成ではない、という位置づけです。
既に動いているジョブがある場合の切り替え手順
もし今すでに手動で起動したジョブが動いている状態で、上の service を start すると、二重起動する可能性があります。
なので、いったん全部止めてから切り替えるのが安全です。
1. いま動いているジョブを確認
cd /
sudo -u flink /opt/flink/current/bin/flink list -r
2. 重複ジョブを全部停止
表示された JobID を全部 cancel します。
cd /
sudo -u flink /opt/flink/current/bin/flink cancel <JOB_ID>
何も出なくなるまで確認:
cd /
sudo -u flink /opt/flink/current/bin/flink list -r
systemd に反映する
sudo systemctl daemon-reload
sudo systemctl enable flink-job-syslog-events.service
sudo systemctl enable flink-job-authlog-events.service
次にジョブを投入。
sudo systemctl start flink-job-syslog-events.service
sudo systemctl start flink-job-authlog-events.service
注意: cancel → 再投入だけだと状態を安全に引き継げない場合がある
本記事では、既存ジョブをいったん停止して再投入する手順を紹介しています。
ただし、状態を持つ Flink ジョブ(集計・upsert・ウィンドウ処理など) を
運用している場合は、単純な cancel → 再投入だけでは状態を安全に引き継げないことがあります。特に Kafka Source で latest-offset を使っている場合、
停止中に溜まったメッセージを読み飛ばすリスクがあります。
動作確認
service 状態確認
sudo systemctl status flink --no-pager
sudo systemctl status flink-job-syslog-events --no-pager -l
sudo systemctl status flink-job-authlog-events --no-pager -l
Flink 側確認
cd /
sudo -u flink /opt/flink/current/bin/flink list -r
想定されるのは この 2 本だけ です。
insert-into_hive_prod.logs.syslog_eventsinsert-into_hive_prod.logs.authlog_events
これ以上増えていたら異常です。
oneshot service の注意点(重要)
今回ハマりやすかったポイントです。
start しても再実行されないことがある
この service は Type=oneshot + RemainAfterExit=yes なので、一度成功すると:
Active: active (exited)
になります。
この状態で
sudo systemctl start flink-job-syslog-events.service
を打っても、すでに active なので再実行されません。
つまり、手動再試験したいときは:
sudo systemctl restart --no-block flink-job-syslog-events.service
sudo systemctl restart --no-block flink-job-authlog-events.service
または:
sudo systemctl stop flink-job-syslog-events.service
sudo systemctl start --no-block flink-job-syslog-events.service
のようにします。
ログ確認
service ログ
journalctl -u flink-job-syslog-events -n 100 --no-pager
journalctl -u flink-job-authlog-events -n 100 --no-pager
ラッパースクリプトログ
sudo cat /var/log/flink/insert-into_hive_prod.logs.syslog_events.log
sudo cat /var/log/flink/insert-into_hive_prod.logs.authlog_events.log
よくある確認ポイント
1. already running. skip submit.
→ すでに同名ジョブが動いているだけなので正常
2. sql-client submit failed
→ SQL / Kafka connector / Iceberg catalog / sink テーブル定義を確認
3. submitted successfully なのに flink list -r に出ない
→ submit は通ったが、その後 Flink ジョブ本体が起動失敗している可能性あり
その場合は Flink 本体ログも見ます。
sudo ls -ltr /opt/flink/current/log
sudo ls -ltr /opt/flink/current/logs
存在するほうで:
sudo tail -n 200 /opt/flink/current/log/*standalonesession*.log
sudo tail -n 200 /opt/flink/current/log/*taskexecutor*.log
再起動試験
最後に、OS 再起動しても自動復帰するか確認します。
sudo reboot
再起動後:
sudo systemctl status flink --no-pager
sudo systemctl status flink-job-syslog-events --no-pager -l
sudo systemctl status flink-job-authlog-events --no-pager -l
cd /
sudo -u flink /opt/flink/current/bin/flink list -r
ここで 2 本だけ動いていれば完成です。
補足: find: Failed to restore initial working directory について
たとえばこういう警告が出ることがあります。
find: Failed to restore initial working directory: /home/naritomo: Permission denied
これは、sudo -u flink で実行したときに、元のカレントディレクトリ /home/naritomo に flink ユーザーが戻れないだけです。
本質的な問題ではありません。
回避したいなら、実行前に / へ移動します。
cd /
sudo -u flink /opt/flink/current/bin/flink list -r
今回の service では WorkingDirectory=/ を入れているので、この警告は出にくくなります。
今回のポイントまとめ
今回の学びはこの 4 つです。
1. sql-client.sh を Restart=always で管理してはいけない
submit 型ジョブは systemd 常駐管理と相性が悪いです。
2. Type=oneshot で「1回だけ submit」にする
systemd 側は submit トリガーとして使うのが安全です。
3. スクリプト側でも「既に同名ジョブがあれば submit しない」を入れる
これで重複起動事故をかなり防げます。
4. oneshot service は start では再実行されないことがある
手動再試験時は restart を使うのが安全です。
おわりに
これで、
- Flink クラスタは自動起動
- SQL ジョブも自動投入
- 同名ジョブがあれば skip
- ジョブ増殖を防止
- 手動再試験も正しくできる
という構成になりました。
Flink SQL を systemd に載せると、最初は「とりあえず service 化すればいいだろう」とやりがちですが、
submit 型ジョブを常駐 service と同じノリで扱うと事故る、というのが今回の落とし穴でした。
同じように、
- Kafka → Iceberg 継続投入
- Flink SQL の運用自動化
- サーバ再起動時の自動復旧
をやりたい方の参考になれば幸いです。
補足: 今回の構成は「まず動かす」ことを優先した実践例です
本記事では、自宅ラボ / 検証環境で Flink + Kafka + Iceberg を実際に動かして理解することを優先して構成しています。
そのため、実運用に寄せる場合は以下の観点を追加で検討するとより安全です。
savepoint / checkpoint を使った状態復元
Kafka オフセット継続の考慮
ジョブ死活監視(systemd だけに依存しない)
時刻 / タイムゾーンの厳密な扱い
watermark / 遅延到着データの設計「まず動く」から「止めても安全・更新しても安全」へ進めると、
より本番に近い構成になります。