0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

Flink SQL ジョブを systemd で自動起動させた

0
Last updated at Posted at 2026-04-05

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
    • syslog
    • authlog
  • 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_events
  • insert-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/naritomoflink ユーザーが戻れないだけです。

本質的な問題ではありません。

回避したいなら、実行前に / へ移動します。

cd /
sudo -u flink /opt/flink/current/bin/flink list -r

今回の service では WorkingDirectory=/ を入れているので、この警告は出にくくなります。


今回のポイントまとめ

今回の学びはこの 4 つです。

1. sql-client.shRestart=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 / 遅延到着データの設計

「まず動く」から「止めても安全・更新しても安全」へ進めると、
より本番に近い構成になります。

0
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?