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?

Apache Flink + Iceberg で syslog/authlog の1分集計テーブルを作り、Trino/Grafana を高速化する

0
Last updated at Posted at 2026-04-05

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_1m
  • hive_prod.logs.authlog_host_1m

保存するのは以下の最小限の情報だけです。

  • window_start … 1分窓の開始時刻
  • host
  • cnt … その1分間の件数
  • dt … 日付(パーティション用)

この構成にすると、Grafana では 集計済みの件数テーブルを見るだけでよくなるため、かなり軽くなります。


この構成のメリット

この構成にするメリットは次のとおりです。

  • Grafana の分単位グラフが軽くなる
  • Trino の GROUP BY minute, host 負荷を減らせる
  • イベント明細と可視化用集計の責務を分離できる
  • 1日保持にすれば運用がシンプル

つまり、役割を次のように分けられます。

イベント明細テーブル

  • syslog_events
  • authlog_events

用途:

  • 個別ログの調査
  • メッセージ本文確認
  • 障害解析
  • フィルタ検索

1分集計テーブル

  • syslog_host_1m
  • authlog_host_1m

用途:

  • Grafana 可視化
  • 時系列件数グラフ
  • ホスト別ランキング
  • ダッシュボード用

前提

本記事は、以下がすでに構築済みである前提です。

  • Kafka
  • Flink
  • Hive Metastore
  • Iceberg
  • HDFS
  • Trino
  • Grafana

以下の記事の続き・派生記事という位置づけです。


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_start
  • host
  • dt

に対しては、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本です。

  • syslogsyslog_host_1m
  • authlogauthlog_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_starttimestamp として保持しているため、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_events
  • authlog_events
  • syslog_host_1m
  • authlog_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 / 遅延到着データの設計

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

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?