Apache Flink + Iceberg で syslog/authlog の1分集計テーブルを作り、Trino/Grafana を高速化する
はじめに
これまで、Kafka → Flink → Iceberg で syslog / authlog のイベント明細テーブルを作成し、Trino や Grafana から可視化できるようにしてきました。
ただし、Grafana で以下のような可視化をしようとすると、イベント明細テーブルをそのまま Trino で毎回集計するのは重くなりがちです。
- 分毎・ホスト毎の syslog 件数
- 分毎・ホスト毎の authlog 件数
- ホスト毎の syslog 件数
- ホスト毎の authlog 件数
例えば、以下のようなクエリを毎回イベント明細に対して実行すると、件数が増えるにつれて重くなっていきます。
SELECT
CAST(date_trunc('minute', ts) AS timestamp) AS time,
host AS metric,
COUNT(*) AS value
FROM iceberg.logs.syslog_events
WHERE ts BETWEEN CAST(from_iso8601_timestamp('${__from:date:iso}') AS timestamp)
AND CAST(from_iso8601_timestamp('${__to:date:iso}') AS timestamp)
GROUP BY 1, 2
ORDER BY 1, 2;
そこで今回は、Flink で1分単位に事前集計した専用テーブルを Iceberg に持ち、Trino/Grafana 側のクエリを軽くする構成を作ります。
今回やること
以下の 1分集計専用テーブル を追加します。
hive_prod.logs.syslog_host_1mhive_prod.logs.authlog_host_1m
保存するのは以下の最小限の情報だけです。
-
window_start… 1分窓の開始時刻 host-
cnt… その1分間の件数 -
dt… 日付(パーティション用)
この構成にすると、Grafana では 集計済みの件数テーブルを見るだけでよくなるため、かなり軽くなります。
この構成のメリット
この構成にするメリットは次のとおりです。
- Grafana の分単位グラフが軽くなる
- Trino の
GROUP BY minute, host負荷を減らせる - イベント明細と可視化用集計の責務を分離できる
- 1日保持にすれば運用がシンプル
つまり、役割を次のように分けられます。
イベント明細テーブル
syslog_eventsauthlog_events
用途:
- 個別ログの調査
- メッセージ本文確認
- 障害解析
- フィルタ検索
1分集計テーブル
syslog_host_1mauthlog_host_1m
用途:
- Grafana 可視化
- 時系列件数グラフ
- ホスト別ランキング
- ダッシュボード用
前提
本記事は、以下がすでに構築済みである前提です。
- Kafka
- Flink
- Hive Metastore
- Iceberg
- HDFS
- Trino
- Grafana
以下の記事の続き・派生記事という位置づけです。
-
Kafka → Flink → Iceberg でイベント明細テーブルを作る手順
-
Flink ジョブを systemd で自動起動する手順
1. Iceberg の1分集計テーブルを作成する
今回は 1日保持 で使う想定なので、パーティションは dt のみとします。
ope1から行うこと。(HDFS保管先作成)
sudo -u hadoop hdfs dfs -mkdir -p /warehouse/iceberg/logs/syslog_host_1m
sudo -u hadoop hdfs dfs -mkdir -p /warehouse/iceberg/logs/authlog_host_1m
flink1から行うこと。(Icebergテーブル作成)
cd /
sudo tee /opt/flink/jobs/flink_create_1m_tables.sql >/dev/null <<'EOF'
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;
CREATE DATABASE IF NOT EXISTS logs;
USE logs;
CREATE TABLE syslog_host_1m (
window_start TIMESTAMP(3),
host STRING,
cnt BIGINT,
dt DATE,
hr INT,
PRIMARY KEY (window_start, host, dt) NOT ENFORCED
)
PARTITIONED BY (dt)
WITH (
'format-version' = '2',
'write.upsert.enabled' = 'true',
'write.distribution-mode' = 'hash',
'location' = 'hdfs://cluster1/warehouse/iceberg/logs/syslog_host_1m'
);
CREATE TABLE authlog_host_1m (
window_start TIMESTAMP(3),
host STRING,
cnt BIGINT,
dt DATE,
hr INT,
PRIMARY KEY (window_start, host, dt) NOT ENFORCED
)
PARTITIONED BY (dt)
WITH (
'format-version' = '2',
'write.upsert.enabled' = 'true',
'write.distribution-mode' = 'hash',
'location' = 'hdfs://cluster1/warehouse/iceberg/logs/authlog_host_1m'
);
EOF
sudo -u flink /opt/flink/current/bin/sql-client.sh embedded -f /opt/flink/jobs/flink_create_1m_tables.sql
確認:
zeppelinで確認する。
%spark.sql
SHOW TABLES IN hive_prod.logs;
DESCRIBE TABLE hive_prod.logs.syslog_host_1m;
DESCRIBE TABLE hive_prod.logs.authlog_host_1m;
補足: location は明示しておくのが安全
今回の構成では、Flink Iceberg DDL で location を明示している。
WITH (
'format-version' = '2',
'write.upsert.enabled' = 'true',
'write.distribution-mode' = 'hash',
'location' = 'hdfs://cluster1/warehouse/iceberg/logs/syslog_host_1m'
)
これは、Hive Metastore / default warehouse 側 (/user/hive/warehouse/...) と
混在してしまう事故を避けるためである。
特に、過去に default warehouse 側でテーブル作成・削除を行っていた環境では、
- Spark では見える
- Trino では metadata file access error
- 再作成後も旧 LOCATION を見に行く
といったトラブルが起きやすい。
そのため、本記事では Iceberg テーブルの実体パスを明示的に固定している。
補足: 集計テーブルに PK を付けている理由
今回の syslog_host_1m / authlog_host_1m は、
「1分 × host ごとの件数」 を保持する集計テーブルである。
そのため、同じキー:
window_starthostdt
に対しては、1行だけ存在する状態が望ましい。
そこで、以下のように PK を定義している。
PRIMARY KEY (window_start, host, dt) NOT ENFORCED
これは RDB の厳密な一意制約ではなく、
Flink / Iceberg に対して「この粒度の行は更新対象として扱う」ためのキーである。
これにより、ジョブ再起動・再投入・再計算時に、
同じ minute / host の行が単純追記されるのを抑えやすくなる。
逆に、
syslog_events/authlog_eventsのような 生イベントテーブル には、
自然な一意キーが定めにくいため、基本的には PK を付けない方が扱いやすい。
2. Flink ジョブ定義
今回は、Kafka topic から直接1分集計して Iceberg に保存します。
作成するジョブは以下の2本です。
-
syslog→syslog_host_1m -
authlog→authlog_host_1m
本作業はFlinkで行うこと
2.1. Flink SQL ファイルを配置する
sudo tee /opt/flink/jobs/flink_syslog_host_1m.sql >/dev/null <<'EOF'
SET 'pipeline.name' = 'insert-into_hive_prod.logs.syslog_host_1m';
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_1m (
ts STRING,
src_host STRING,
program STRING,
message STRING,
event_time AS TO_TIMESTAMP(REPLACE(SUBSTRING(ts, 1, 19), 'T', ' ')),
WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'syslog',
'properties.bootstrap.servers' = 'kafka1:9092,kafka2:9092,kafka3:9092',
'properties.group.id' = 'flink-syslog-host-1m',
'scan.startup.mode' = 'latest-offset',
'format' = 'json',
'json.ignore-parse-errors' = 'true'
);
INSERT INTO hive_prod.logs.syslog_host_1m /*+ OPTIONS(
'compression-codec'='zstd',
'upsert-enabled'='true',
'distribution-mode'='hash',
'write-parallelism'='1',
'target-file-size-bytes'='134217728'
) */
SELECT
window_start,
host,
cnt,
CAST(window_start AS DATE) AS dt,
CAST(EXTRACT(HOUR FROM window_start) AS INT) AS hr
FROM (
SELECT
window_start,
TRIM(src_host) AS host,
COUNT(*) AS cnt
FROM TABLE(
TUMBLE(
TABLE syslog_json_1m,
DESCRIPTOR(event_time),
INTERVAL '1' MINUTE
)
)
WHERE src_host IS NOT NULL
AND TRIM(src_host) <> ''
GROUP BY window_start, TRIM(src_host)
);
EOF
sudo tee /opt/flink/jobs/flink_authlog_host_1m.sql >/dev/null <<'EOF'
SET 'pipeline.name' = 'insert-into_hive_prod.logs.authlog_host_1m';
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_1m (
ts STRING,
src_host STRING,
program STRING,
message STRING,
event_time AS TO_TIMESTAMP(REPLACE(SUBSTRING(ts, 1, 19), 'T', ' ')),
WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'authlog',
'properties.bootstrap.servers' = 'kafka1:9092,kafka2:9092,kafka3:9092',
'properties.group.id' = 'flink-authlog-host-1m',
'scan.startup.mode' = 'latest-offset',
'format' = 'json',
'json.ignore-parse-errors' = 'true'
);
INSERT INTO hive_prod.logs.authlog_host_1m /*+ OPTIONS(
'compression-codec'='zstd',
'upsert-enabled'='true',
'distribution-mode'='hash',
'write-parallelism'='1',
'target-file-size-bytes'='134217728'
) */
SELECT
window_start,
host,
cnt,
CAST(window_start AS DATE) AS dt,
CAST(EXTRACT(HOUR FROM window_start) AS INT) AS hr
FROM (
SELECT
window_start,
TRIM(src_host) AS host,
COUNT(*) AS cnt
FROM TABLE(
TUMBLE(
TABLE authlog_json_1m,
DESCRIPTOR(event_time),
INTERVAL '1' MINUTE
)
)
WHERE src_host IS NOT NULL
AND TRIM(src_host) <> ''
GROUP BY window_start, TRIM(src_host)
);
EOF
sudo chown flink:flink /opt/flink/jobs/flink_syslog_host_1m.sql /opt/flink/jobs/flink_authlog_host_1m.sql
sudo chmod 644 /opt/flink/jobs/flink_syslog_host_1m.sql /opt/flink/jobs/flink_authlog_host_1m.sql
補足:同じ PK のデータが来た場合どうなるか
今回のテーブルでは upsert を有効にしているため、
'write.upsert.enabled' = 'true'
同じ主キー
(window_start, host, dt)
を持つ新しいデータが来た場合、
別レコードとして増えるのではなく、同じキーの行を更新する方向で扱われます。イメージとしては以下です。
例
最初にこう入る:
window_start=2026-04-05 11:23:00
host=master1
dt=2026-04-05
cnt=10その後、同じ PK でこう入る:
window_start=2026-04-05 11:23:00
host=master1
dt=2026-04-05
cnt=15この場合、「11:23 の master1 の集計結果」が 10 → 15 に更新されるイメージです。
補足: 遅延到着データ(late event)について
今回の例では watermark を設定しているため、
ある程度の遅延データを考慮したウィンドウ集計ができます。ただし、どれだけ遅れて届くデータを許容するかによって、
集計の確定タイミング
再更新の発生頻度
upsert の更新回数が変わります。
そのため本番寄りにする場合は、
ログ配送の遅延がどれくらいあるか
数十秒で十分か、数分必要かを見ながら watermark の設定を調整するのがおすすめです。
3. systemd で自動起動する
Flink SQL ジョブは 常駐 service にするのではなく、submit 専用の oneshot service として扱うのが安全です。
本作業はFlinkで行うこと
3-1. syslog_host_1m 用 service
ファイル:
/etc/systemd/system/flink-job-syslog-host-1m.service
sudo tee /etc/systemd/system/flink-job-syslog-host-1m.service >/dev/null <<EOF
[Unit]
Description=Submit Flink SQL Job - syslog_host_1m
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_host_1m /opt/flink/jobs/flink_syslog_host_1m.sql
RemainAfterExit=yes
TimeoutStartSec=300
[Install]
WantedBy=multi-user.target
EOF
3-2. authlog_host_1m 用 service
ファイル:
/etc/systemd/system/flink-job-authlog-host-1m.service
sudo tee /etc/systemd/system/flink-job-authlog-host-1m.service >/dev/null <<EOF
[Unit]
Description=Submit Flink SQL Job - authlog_host_1m
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_host_1m /opt/flink/jobs/flink_authlog_host_1m.sql
RemainAfterExit=yes
TimeoutStartSec=300
[Install]
WantedBy=multi-user.target
EOF
3-3. 有効化と起動
sudo systemctl daemon-reload
sudo systemctl enable flink-job-syslog-host-1m.service
sudo systemctl enable flink-job-authlog-host-1m.service
sudo systemctl start flink-job-syslog-host-1m.service
sudo systemctl start flink-job-authlog-host-1m.service
3-4. 状態確認
sudo systemctl status flink --no-pager
sudo systemctl status flink-job-syslog-host-1m.service --no-pager -l
sudo systemctl status flink-job-authlog-host-1m.service --no-pager -l
sudo -u flink /opt/flink/current/bin/flink list -r
journalctl -u flink-job-syslog-host-1m -n 100 --no-pager
journalctl -u flink-job-authlog-host-1m -n 100 --no-pager
sudo cat /var/log/flink/insert-into_hive_prod.logs.syslog_host_1m.log
sudo cat /var/log/flink/insert-into_hive_prod.logs.authlog_host_1m.log
補足
oneshot + RemainAfterExit=yes なので、正常終了後は以下のようになります。
active (exited)
これは異常ではなく、**「submit が成功して systemd の役目は終わった」**という状態です。
4. データが入っているか確認する
Spark SQL で確認します。
4-1. Spark SQL で確認
zeppelinのspark.sqlから作成してください。
%spark.sql
SELECT *
FROM hive_prod.logs.syslog_host_1m
ORDER BY window_start DESC
LIMIT 20;
%spark.sql
SELECT *
FROM hive_prod.logs.authlog_host_1m
ORDER BY window_start DESC
LIMIT 20;
5. Trino / Grafana 用クエリ
trino1にて、以下コマンドを入力し、参照できることを確認する。
trino --server http://localhost:8080 --execute "
SELECT *
FROM iceberg.logs.syslog_host_1m
ORDER BY window_start DESC
LIMIT 10;
"
trino --server http://localhost:8080 --execute "
SELECT *
FROM iceberg.logs.authlog_host_1m
ORDER BY window_start DESC
LIMIT 10;
"
Grafanaで以下クエリのグラフを作成してください
5-1. 時系列 syslog 件数(分毎・ホスト毎)
SELECT
CAST(date_trunc('minute', window_start) AS timestamp) AS time,
host AS metric,
SUM(cnt) AS value
FROM iceberg.logs.syslog_host_1m
WHERE window_start BETWEEN CAST(from_iso8601_timestamp('${__from:date:iso}') AS timestamp)
AND CAST(from_iso8601_timestamp('${__to:date:iso}') AS timestamp)
GROUP BY 1,2
ORDER BY 1,2
時単位のを取りたい場合、以下のようにすればよい。
SELECT
CAST(date_trunc('hour', window_start) AS timestamp) AS time,
host AS metric,
SUM(cnt) AS value
FROM iceberg.logs.syslog_host_1m
WHERE window_start BETWEEN CAST(from_iso8601_timestamp('${__from:date:iso}') AS timestamp)
AND CAST(from_iso8601_timestamp('${__to:date:iso}') AS timestamp)
GROUP BY 1,2
ORDER BY 1,2
5-2. 時系列 authlog 件数(分毎・ホスト毎)
SELECT
CAST(date_trunc('minute', window_start) AS timestamp) AS time,
host AS metric,
SUM(cnt) AS value
FROM iceberg.logs.authlog_host_1m
WHERE window_start BETWEEN CAST(from_iso8601_timestamp('${__from:date:iso}') AS timestamp)
AND CAST(from_iso8601_timestamp('${__to:date:iso}') AS timestamp)
GROUP BY 1,2
ORDER BY 1,2
時単位のを取りたい場合、以下のようにすればよい。
SELECT
CAST(date_trunc('hour', window_start) AS timestamp) AS time,
host AS metric,
SUM(cnt) AS value
FROM iceberg.logs.authlog_host_1m
WHERE window_start BETWEEN CAST(from_iso8601_timestamp('${__from:date:iso}') AS timestamp)
AND CAST(from_iso8601_timestamp('${__to:date:iso}') AS timestamp)
GROUP BY 1,2
ORDER BY 1,2
5-3. ホスト毎の syslog 件数
SELECT
host,
SUM(cnt) AS total_cnt
FROM iceberg.logs.syslog_host_1m
WHERE window_start BETWEEN CAST(from_iso8601_timestamp('${__from:date:iso}') AS timestamp)
AND CAST(from_iso8601_timestamp('${__to:date:iso}') AS timestamp)
GROUP BY host
ORDER BY total_cnt DESC
5-4. ホスト毎の authlog 件数
SELECT
host,
SUM(cnt) AS total_cnt
FROM iceberg.logs.authlog_host_1m
WHERE window_start BETWEEN CAST(from_iso8601_timestamp('${__from:date:iso}') AS timestamp)
AND CAST(from_iso8601_timestamp('${__to:date:iso}') AS timestamp)
GROUP BY host
ORDER BY total_cnt DESC
補足
-
window_startで実時間帯を絞る -
window_startはtimestampとして保持しているため、Grafana の${__from}/${__to}も
timestampにキャストして比較している
これにより、表示側で余計なタイムゾーン変換を重ねずに比較できる。
6. 各種運用シェル修正
6-1. データ保持は1日で十分
今回の 1分集計テーブルは、Grafana 可視化用の軽量テーブルです。
長期保持の役割はイベント明細テーブルに任せ、こちらは 1日分だけ保持 で十分です。
古いデータ削除については、既存の cleanup /orphan files Cleanシェルに組み込んで運用します。
6-2. 日次 cleanup シェル修正
ope1 に配置します。
/opt/iceberg/bin/cleanup_flink_aggregate_daily.sh
TABLES=() 内のテーブルを追加する。
sudo tee /opt/iceberg/bin/cleanup_flink_aggregate_daily.sh >/dev/null <<'EOF'
#!/usr/bin/env bash
set -euo pipefail
SPARK_SQL="${SPARK_SQL:-sudo -u spark /usr/local/bin/spark-sql-iceberg}"
CATALOG="${CATALOG:-hive_prod}"
DB="${DB:-logs}"
RETAIN_LAST="${RETAIN_LAST:-1}"
CUTOFF_TS="${CUTOFF_TS:-$(date -d '1 day ago 00:00:00' '+%F %T')}"
LOG_DIR="${LOG_DIR:-/tmp}"
LOG_FILE="${LOG_DIR}/cleanup_flink_aggregate_daily.log"
mkdir -p "${LOG_DIR}"
TABLES=(
"authlog_events"
"syslog_events"
"syslog_host_1m"
"authlog_host_1m"
)
log() {
local msg="[$(date '+%F %T')] $*"
echo "${msg}"
echo "${msg}" >> "${LOG_FILE}" 2>/dev/null || true
}
run_sql() {
local sql="$1"
bash -lc "${SPARK_SQL} <<SQL
${sql}
SQL
"
}
log "START cleanup flink aggregate daily cutoff=${CUTOFF_TS}"
for tbl in "${TABLES[@]}"; do
FULL_TABLE="${CATALOG}.${DB}.${tbl}"
log "DELETE old rows from ${FULL_TABLE}"
run_sql "
DELETE FROM ${FULL_TABLE}
WHERE dt < current_date();
" | tee -a "${LOG_FILE}"
log "EXPIRE SNAPSHOTS for ${FULL_TABLE}"
run_sql "
CALL ${CATALOG}.system.expire_snapshots(
table => '${DB}.${tbl}',
older_than => TIMESTAMP '${CUTOFF_TS}',
retain_last => ${RETAIN_LAST},
clean_expired_metadata => true
);
" | tee -a "${LOG_FILE}"
done
log "END cleanup flink aggregate daily"
EOF
6-3. orphan files クリア用シェル修正
ope1 に配置します。
/opt/iceberg/bin/cleanup_flink_aggregate_orphan.sh
TABLES=() 内のテーブルを追加する。
sudo tee /opt/iceberg/bin/cleanup_flink_aggregate_orphan.sh >/dev/null <<'EOF'
#!/usr/bin/env bash
set -euo pipefail
SPARK_SQL="${SPARK_SQL:-sudo -u spark /usr/local/bin/spark-sql-iceberg}"
CATALOG="${CATALOG:-hive_prod}"
DB="${DB:-logs}"
LOG_DIR="${LOG_DIR:-/tmp}"
LOG_FILE="${LOG_DIR}/cleanup_flink_aggregate_orphan.log"
# true: 候補確認のみ
# false: 実際に削除
DRY_RUN="${DRY_RUN:-true}"
# 実削除時のしきい値
ORPHAN_CUTOFF_TS="${ORPHAN_CUTOFF_TS:-$(date -d '24 hours ago' '+%F %T')}"
mkdir -p "${LOG_DIR}"
TABLES=(
"authlog_events"
"syslog_events"
"syslog_host_1m"
"authlog_host_1m"
)
log() {
local msg="[$(date '+%F %T')] $*"
echo "${msg}"
echo "${msg}" >> "${LOG_FILE}" 2>/dev/null || true
}
run_sql() {
local sql="$1"
bash -lc "${SPARK_SQL} <<SQL
${sql}
SQL
"
}
log "START orphan cleanup cutoff=${ORPHAN_CUTOFF_TS}"
for tbl in "${TABLES[@]}"; do
FULL_TABLE="${CATALOG}.${DB}.${tbl}"
if [[ "${DRY_RUN}" == "true" ]]; then
log "REMOVE ORPHAN FILES DRY-RUN for ${FULL_TABLE}"
run_sql "
CALL ${CATALOG}.system.remove_orphan_files(
table => '${DB}.${tbl}',
dry_run => true
);
" | tee -a "${LOG_FILE}"
else
log "REMOVE ORPHAN FILES for ${FULL_TABLE}"
run_sql "
CALL ${CATALOG}.system.remove_orphan_files(
table => '${DB}.${tbl}',
older_than => TIMESTAMP '${ORPHAN_CUTOFF_TS}'
);
" | tee -a "${LOG_FILE}"
fi
done
log "END orphan cleanup"
EOF
補足: 保持対象の考え方
今回の cleanup 例では、以下の4テーブルを対象にしている。
syslog_eventsauthlog_eventssyslog_host_1mauthlog_host_1m
つまり、1分集計テーブルだけでなく、イベント明細テーブルも短期保持する前提である。
これは、
- Grafana 可視化
- 直近の調査
- 軽量運用
を優先した構成である。
もし、
- 生イベントはもう少し長く保持したい
- 1分集計だけ短期ローテーションしたい
という場合は、cleanup 対象を分けてもよい。
例:
-
*_host_1m→ 1日保持 -
*_events→ 数日〜数週間保持
6-4. compactionスクリプト修正
ope1 に配置します。
/opt/iceberg/bin/compact_iceberg.sh
TABLES=() 内のテーブルを追加する。
tee /opt/iceberg/bin/compact_iceberg.sh <<'EOF'
#!/usr/bin/env bash
set -euo pipefail
SPARK_SQL="${SPARK_SQL:-sudo -u spark /usr/local/bin/spark-sql-iceberg}"
CATALOG="${CATALOG:-hive_prod}"
DB="${DB:-logs}"
LOG_DIR="${LOG_DIR:-/tmp}"
LOG_FILE="${LOG_DIR}/compact_iceberg.log"
TABLES=(
"syslog_iceberg"
"authlog_iceberg"
"authlog_events"
"syslog_events"
"syslog_host_1m"
"authlog_host_1m"
)
mkdir -p "${LOG_DIR}"
log() {
local msg="[INFO] $(date '+%F %T') $*"
echo "${msg}"
echo "${msg}" >> "${LOG_FILE}" 2>/dev/null || true
}
err() {
local msg="[ERROR] $(date '+%F %T') $*"
echo "${msg}" >&2
echo "${msg}" >> "${LOG_FILE}" 2>/dev/null || true
}
run_sql() {
local sql="$1"
bash -lc "${SPARK_SQL} <<SQL
${sql}
SQL
"
}
log "START compact iceberg tables"
for tbl in "${TABLES[@]}"; do
log "rewrite_data_files start: ${CATALOG}.${DB}.${tbl}"
if run_sql "
CALL ${CATALOG}.system.rewrite_data_files(
table => '${DB}.${tbl}'
);
" >> "${LOG_FILE}" 2>&1; then
log "rewrite_data_files done : ${CATALOG}.${DB}.${tbl}"
else
err "rewrite_data_files failed: ${CATALOG}.${DB}.${tbl}"
fi
done
log "END compact iceberg tables"
EOF
6-5. iceberg HDFS空フォルダ掃除
icebergに対し、運用していくとデータフォルダに空フォルダが蓄積されてしまうため、
下記シェルを作成し、上記処理の最後に追加する。
ope1 に配置します。
/opt/iceberg/bin/cleanup_empty_hdfs_dirs_recursive.sh
tee /opt/iceberg/bin/cleanup_empty_hdfs_dirs_recursive.sh >/dev/null <<'EOF'
#!/usr/bin/env bash
set -euo pipefail
echo "[INFO] $(date '+%F %T') cleanup start"
CUTOFF_14="$(date -d '14 days ago' +%F)"
echo "[INFO] cutoff_14=${CUTOFF_14}"
sudo -u hadoop bash -s -- "$CUTOFF_14" <<'EOS'
set -euo pipefail
CUTOFF_14="$1"
for base in \
/warehouse/iceberg/logs/authlog_events/data \
/warehouse/iceberg/logs/authlog_host_1m/data \
/warehouse/iceberg/logs/syslog_events/data \
/warehouse/iceberg/logs/syslog_host_1m/data \
/warehouse/iceberg/logs/authlog_iceberg/data \
/warehouse/iceberg/logs/syslog_iceberg/data
do
echo "[INFO] base=$base"
hdfs dfs -ls "$base" 2>/dev/null | awk '$1 ~ /^d/ {print $8}' | while read -r dt_dir; do
echo "[INFO] dt_dir=$dt_dir"
dir_name="$(basename "$dt_dir")"
# authlog_iceberg / syslog_iceberg は直近14日を対象外にする
if [ "$base" = "/warehouse/iceberg/logs/authlog_iceberg/data" ] || \
[ "$base" = "/warehouse/iceberg/logs/syslog_iceberg/data" ]; then
case "$dir_name" in
dt=????-??-??)
part_date="${dir_name#dt=}"
if [ "$part_date" \> "$CUTOFF_14" ] || [ "$part_date" = "$CUTOFF_14" ]; then
echo "[INFO] skip recent iceberg dir=$dt_dir"
continue
fi
;;
esac
fi
hdfs dfs -ls "$dt_dir" 2>/dev/null | awk '$1 ~ /^d/ {print $8}' | while read -r subdir; do
hdfs dfs -rmdir "$subdir" 2>/dev/null || true
done
hdfs dfs -rmdir "$dt_dir" 2>/dev/null || true
done
done
EOS
echo "[INFO] $(date '+%F %T') cleanup done"
EOF
権限付与
chmod 755 /opt/iceberg/bin/cleanup_empty_hdfs_dirs_recursive.sh
手動実行確認
/opt/iceberg/bin/cleanup_empty_hdfs_dirs_recursive.sh
cron 登録
毎日 01:30 に実行する例です。
30 1 * * * /opt/iceberg/bin/cleanup_empty_hdfs_dirs_recursive.sh >/dev/null 2>&1
7. まとめ
今回、Flink を使って Kafka 上の syslog / authlog を 1分単位 × host 単位 で集計し、
Iceberg に格納して Trino / Grafana から軽量に可視化できるようにした。
この構成の利点は以下。
- Trino で minute 集計を毎回生計算しなくてよい
- Grafana から軽く時系列可視化できる
-
syslog/authlogの件数変化をすぐ見られる - host別の負荷・異常傾向も追いやすい
- 1日分だけ保持することで軽量運用しやすい
特に今回のように、
- Kafka でログを受ける
- Iceberg に蓄積する
- Trino / Grafana で見る
という流れでは、**「集計済みの軽いテーブルを別に持つ」**のがかなり効く。
今後は例えば以下にも拡張しやすい。
- 5分 / 1時間集計テーブル
- authlog 成功 / 失敗別集計
- user別集計
- host別 TopN パネル
- Flink checkpoint / savepoint を含めた運用強化
補足: 今回の構成は「まず動かす」ことを優先した実践例です
本記事では、自宅ラボ / 検証環境で Flink + Kafka + Iceberg を実際に動かして理解することを優先して構成しています。
そのため、実運用に寄せる場合は以下の観点を追加で検討するとより安全です。
savepoint / checkpoint を使った状態復元
Kafka オフセット継続の考慮
ジョブ死活監視(systemd だけに依存しない)
時刻 / タイムゾーンの厳密な扱い
watermark / 遅延到着データの設計「まず動く」から「止めても安全・更新しても安全」へ進めると、
より本番に近い構成になります。