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?

ログ収集・解析基盤にSpark/Icebergを構築してみた

0
Last updated at Posted at 2026-03-22

はじめに

完成させたログ収集・解析基盤にSpark/Icebergを構築してみました。

前提として以下の記事を参照し、解析基盤を構築していること。

追記: AWS S3 / Glue / Athena への分析用コピー

本記事では、自宅 HDFS 上に hive_prod.logs.syslog_iceberg / hive_prod.logs.authlog_iceberg を作成し、Spark / Trino から分析できる状態にするところまでを扱っています。

この Iceberg テーブルを Source of Truth として残したまま、AWS 側に分析用コピーを作る構成については、別記事で整理しました。

ポイントは、HDFS 上の Iceberg warehouse を aws s3 sync で単純コピーするのではなく、Spark + Iceberg で読み出し、AWS 側に glue_prod.logs.syslog_iceberg / glue_prod.logs.authlog_iceberg として独立した Iceberg テーブルを作ることです。

Athena から SQL 分析できるようになった後は、QuickSight / Amazon Quick から可視化する構成にも広げられます。

Spark/Iceberg構築

構成

  • 既存基盤
ホスト名 役割
master1/master2/master3 HDFS / YARN / Hive Metastore / HiveServer2
worker1/worker2/... DataNode / NodeManager
  • 新規追加ホスト
ホスト名 役割
iceberg1 Spark Master / Spark History Server
iceberg2 Spark Worker
iceberg3 Spark Worker

1ホストで動かす場合は以下の構成でも動かせます。

ホスト名 役割
iceberg1 Spark Master / Spark History Server / Spark Worker
  • 既存 curated は残す
    例: logs.syslog_curated, authlogs.authlog_curated

  • Iceberg は別テーブルで追加
    例: logs.syslog_iceberg, logs.authlog_iceberg

1. 方針

今回の運用方針

  • 既存 curated は消さない
  • Iceberg はcuratedを基に別テーブルで作る
  • 書き込み・保守は Spark 側
  • Hive 3.1.3 は参照中心
  • 新機能ホストは iceberg1/2/3 に分離

作るもの

  • HDFS 上に Iceberg 専用 warehouse
  • Spark 3.5 系クラスタ
  • Spark から Hive Metastore を使う Iceberg catalog
  • logs.syslog_iceberg
  • logs.authlog_iceberg

2. HDFS 側の準備

ope1 など HDFS 操作可能なホストで実行します。

sudo -u hadoop hdfs dfs -mkdir -p /warehouse/iceberg/logs
sudo -u hadoop hdfs dfs -chown -R spark:hadoop /warehouse/iceberg
sudo -u hadoop hdfs dfs -chmod -R 775 /warehouse/iceberg

sudo -u hadoop hdfs dfs -mkdir -p /user/spark/applicationHistory
sudo -u hadoop hdfs dfs -chown -R spark:hadoop /user/spark/applicationHistory
sudo -u hadoop hdfs dfs -chmod -R 775 /user/spark/applicationHistory

3. 新規ホスト iceberg1/2/3 の OS 前提

全ホスト iceberg1/2/3 で実行:

Ubuntu24.04でホスト構築すること。

sudo apt-get update
sudo apt-get install -y openjdk-11-jdk curl tar
sudo useradd -r -m -d /var/lib/spark -s /bin/bash spark || true
sudo mkdir -p /opt/spark
sudo mkdir -p /etc/spark/conf
sudo chown -R root:root /opt/spark /etc/spark
sudo chmod 755 /opt/spark /etc/spark

確認:

java -version

4. Spark 3.5 系の配置

Spark は 3.5 系を使います。Iceberg は Spark 3.5 向けに専用 runtime artifact を提供しています。

全ホスト iceberg1/2/3 で:

cd /tmp
curl -LO https://archive.apache.org/dist/spark/spark-3.5.8/spark-3.5.8-bin-hadoop3.tgz
sudo tar -xzf spark-3.5.8-bin-hadoop3.tgz -C /opt/spark
sudo ln -sfn /opt/spark/spark-3.5.8-bin-hadoop3 /opt/spark/current
sudo chown -R root:root /opt/spark

確認:

/opt/spark/current/bin/spark-submit --version

5. Hadoop/Hive 設定を Spark 側に配布

iceberg1/2/3 が既存 HDFS と Hive Metastore を使うため、既存クラスタの設定をコピーします。

既存ホストから以下を回収して iceberg1/2/3 へ配置:

/etc/hadoop/conf/core-site.xml
/etc/hadoop/conf/hdfs-site.xml
/etc/hive/conf/hive-site.xml

配置先:

sudo mkdir -p /etc/hadoop/conf
sudo mkdir -p /etc/hive/conf

権限:

sudo chmod 644 /etc/hadoop/conf/*.xml /etc/hive/conf/*.xml

6. Spark 設定

全ホスト iceberg1/2/3 に /etc/spark/conf/spark-env.sh を作成:

sudo tee /etc/spark/conf/spark-env.sh > /dev/null <<'EOF'
export JAVA_HOME=$(dirname $(dirname $(readlink -f $(which java))))
export SPARK_HOME=/opt/spark/current
export SPARK_CONF_DIR=/etc/spark/conf
export HADOOP_CONF_DIR=/etc/hadoop/conf
export HIVE_CONF_DIR=/etc/hive/conf
export SPARK_LOG_DIR=/var/log/spark
export SPARK_PID_DIR=/var/run/spark
export SPARK_WORKER_DIR=/var/lib/spark/work
export SPARK_LOCAL_DIRS=/var/lib/spark/work
export SPARK_MASTER_WEBUI_PORT=18080
export SPARK_WORKER_WEBUI_PORT=18081
EOF
sudo mkdir -p /var/lib/spark/tmp
sudo chown -R spark:spark /var/lib/spark
sudo chmod +x /etc/spark/conf/spark-env.sh

次に /etc/spark/conf/spark-defaults.conf を全ホスト共通で作成します。

sudo tee /etc/spark/conf/spark-defaults.conf > /dev/null <<'EOF'
spark.master spark://iceberg1:7077

spark.sql.extensions org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions

spark.sql.catalog.hive_prod org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.hive_prod.type hive
spark.sql.catalog.hive_prod.uri thrift://hive1:9083,thrift://hive2:9083
spark.sql.catalog.hive_prod.warehouse hdfs://cluster1/warehouse/iceberg

spark.driver.extraClassPath /opt/spark/extra-jars/iceberg-spark-runtime-3.5_2.12-1.10.1.jar
spark.executor.extraClassPath /opt/spark/extra-jars/iceberg-spark-runtime-3.5_2.12-1.10.1.jar

spark.sql.sources.partitionOverwriteMode dynamic
spark.sql.session.timeZone Asia/Tokyo

spark.driver.memory 1g
spark.executor.memory 512m
spark.executor.cores 1
spark.cores.max 1
EOF
sudo mkdir -p /var/log/spark
sudo chown -R spark:spark /var/log/spark
sudo chmod 755 /var/log/spark

sudo mkdir -p /var/run/spark
sudo chown -R spark:spark /var/run/spark
sudo chmod 755 /var/run/spark

sudo mkdir -p /var/lib/spark/work
sudo chown -R spark:spark /var/lib/spark
sudo chmod -R 755 /var/lib/spark

7. Iceberg runtime jar の導入

全ホスト iceberg1/2/3 で:

sudo mkdir -p /opt/spark/extra-jars
cd /opt/spark/extra-jars
sudo curl -LO https://repo1.maven.org/maven2/org/apache/iceberg/iceberg-spark-runtime-3.5_2.12/1.10.1/iceberg-spark-runtime-3.5_2.12-1.10.1.jar
sudo chmod 644 /opt/spark/extra-jars/*.jar

8. systemd ユニット作成

iceberg1: Spark Master

sudo tee /etc/systemd/system/spark-master.service > /dev/null <<'EOF'
[Unit]
Description=Apache Spark Master
After=network.target

[Service]
Type=forking
User=spark
Group=spark
Environment=SPARK_HOME=/opt/spark/current
Environment=SPARK_CONF_DIR=/etc/spark/conf
Environment=SPARK_LOG_DIR=/var/log/spark
Environment=SPARK_PID_DIR=/var/run/spark
ExecStart=/opt/spark/current/sbin/start-master.sh
ExecStop=/opt/spark/current/sbin/stop-master.sh
RemainAfterExit=yes
Environment="JAVA_HOME=/usr/lib/jvm/java-11-openjdk-amd64"
Environment="PATH=/usr/lib/jvm/java-11-openjdk-amd64/bin:/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin"

[Install]
WantedBy=multi-user.target
EOF

iceberg2/3: Spark Worker

sudo tee /etc/systemd/system/spark-worker.service > /dev/null <<'EOF'
[Unit]
Description=Apache Spark Worker
After=network.target

[Service]
Type=forking
User=spark
Group=spark
Environment=SPARK_HOME=/opt/spark/current
Environment=SPARK_CONF_DIR=/etc/spark/conf
Environment=SPARK_LOG_DIR=/var/log/spark
Environment=SPARK_PID_DIR=/var/run/spark
ExecStart=/opt/spark/current/sbin/start-worker.sh spark://iceberg1:7077
ExecStop=/opt/spark/current/sbin/stop-worker.sh
RemainAfterExit=yes
Environment="JAVA_HOME=/usr/lib/jvm/java-11-openjdk-amd64"
Environment="PATH=/usr/lib/jvm/java-11-openjdk-amd64/bin:/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin"

[Install]
WantedBy=multi-user.target
EOF

iceberg1: Spark History Server

sudo tee /etc/systemd/system/spark-history.service > /dev/null <<'EOF'
[Unit]
Description=Apache Spark History Server
After=network.target

[Service]
Type=forking
User=spark
Group=spark
Environment=SPARK_HOME=/opt/spark/current
Environment=SPARK_CONF_DIR=/etc/spark/conf
Environment=SPARK_LOG_DIR=/var/log/spark
Environment=SPARK_PID_DIR=/var/run/spark
ExecStart=/opt/spark/current/sbin/start-history-server.sh
ExecStop=/opt/spark/current/sbin/stop-history-server.sh
RemainAfterExit=yes
Environment="JAVA_HOME=/usr/lib/jvm/java-11-openjdk-amd64"
Environment="PATH=/usr/lib/jvm/java-11-openjdk-amd64/bin:/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin"

[Install]
WantedBy=multi-user.target
EOF

起動:

iceberg1

sudo install -d -o spark -g spark -m 1777 /tmp/spark-events
sudo systemctl daemon-reload
sudo systemctl enable --now spark-master spark-history
sudo systemctl status spark-master --no-pager
sudo systemctl status spark-history --no-pager

iceberg2

sudo systemctl daemon-reload
sudo systemctl enable --now spark-worker
sudo systemctl status spark-worker --no-pager

iceberg3

sudo systemctl daemon-reload
sudo systemctl enable --now spark-worker
sudo systemctl status spark-worker --no-pager

9. Spark から HDFS / Hive Metastore 接続確認

iceberg1 で:

sudo tee /usr/local/bin/spark-sql-iceberg > /dev/null <<'EOF'
#!/usr/bin/env bash
set -euo pipefail

export JAVA_HOME=/usr/lib/jvm/java-11-openjdk-amd64
export PATH="${JAVA_HOME}/bin:/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin"
export SPARK_HOME=/opt/spark/current
export SPARK_CONF_DIR=/etc/spark/conf
export HADOOP_CONF_DIR=/etc/hadoop/conf
export HIVE_CONF_DIR=/etc/hive/conf

exec /opt/spark/current/bin/spark-sql \
  --conf spark.driver.memory=1g \
  --conf spark.executor.memory=512m \
  --conf spark.executor.cores=1 \
  --conf spark.cores.max=1 \
  --conf spark.sql.shuffle.partitions=1 \
  --conf spark.sql.parquet.enableDictionary=false \
  "$@"
EOF

sudo chmod 755 /usr/local/bin/spark-sql-iceberg
sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
SHOW DATABASES;
EOF

これで既存 Hive Metastore の DB 一覧が見えれば、catalog の基本疎通は通っています。

10. 既存 curated は残したまま Iceberg DB を作る

ここからは 既存 curated を消さずに、Iceberg 側を追加します。

iceberg1 で実行:

sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
CREATE DATABASE IF NOT EXISTS hive_prod.logs;
EOF

Spark の Iceberg DDL では CREATE TABLE ... USING iceberg、CREATE TABLE AS SELECT などが使えます。

11. syslog_iceberg 作成

sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
CREATE TABLE hive_prod.logs.syslog_iceberg (
  host       string,
  ts         TIMESTAMP_NTZ,
  severity   int,
  program    string,
  msg        string,
  dt         date,
  hr         int
)
USING iceberg
LOCATION 'hdfs://cluster1/warehouse/iceberg/logs/syslog_iceberg'
PARTITIONED BY (dt, host)
TBLPROPERTIES (
  'format-version'='2',
  'write.distribution-mode'='hash'
);
EOF

12. authlog_iceberg 作成

sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
CREATE TABLE hive_prod.logs.authlog_iceberg (
  host       string,
  ts         TIMESTAMP_NTZ,
  severity   int,
  program    string,
  msg        string,
  dt         date,
  hr         int
)
USING iceberg
LOCATION 'hdfs://cluster1/warehouse/iceberg/logs/authlog_iceberg'
PARTITIONED BY (dt, host)
TBLPROPERTIES (
  'format-version'='2',
  'write.distribution-mode'='hash'
);
EOF

13. 初回ロード

13-0. ts の受け渡しルール

Hive の TIMESTAMP を Parquet 経由で Spark から読む場合、
Hive 側で書いた時の JVM タイムゾーンと Spark の spark.sql.session.timeZone の組み合わせにより、
Iceberg 投入後の ts が 9 時間ずれることがあります。

そのため、Hive curated では次のように役割を分けます。

  • ts ... Hive 内での確認・集計用
  • ts_text ... Spark/Iceberg へ渡すための yyyy-MM-dd HH:mm:ss 文字列
  • ts_raw ... Kafka JSON 由来の元文字列

Spark から Iceberg へ投入するときは、Hive Parquet の TIMESTAMP である ts ではなく、
ts_textTIMESTAMP_NTZ に変換します。

まず ts_rawts_text が同じローカル時刻を表しているか確認します。

sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
SET spark.sql.hive.convertMetastoreParquet=false;
SET spark.sql.session.timeZone=Asia/Tokyo;

SELECT
  host,
  CAST(ts AS STRING) AS ts_seen_by_spark,
  ts_text,
  ts_raw,
  dt
FROM logs.syslog_curated
WHERE dt = date_sub(current_date(), 1)
ORDER BY ts_text DESC
LIMIT 5;
EOF

既存の logs.syslog_curatedts_text / ts_raw がない場合は、
Hive の curated テーブル定義を更新した上で RAW から再生成します。

JDBC_URL='jdbc:hive2://master1:2181,master2:2181,master3:2181/;serviceDiscoveryMode=zooKeeper;zooKeeperNamespace=hiveserver2'

beeline -u "${JDBC_URL}" -n hive -e "
DROP TABLE IF EXISTS logs.syslog_curated;
"

sudo -u hadoop hdfs dfs -rm -r -skipTrash /data/kafka/syslog_curated/* || true

beeline -u "${JDBC_URL}" -n hive -f /usr/local/share/syslog_init.sql
sudo REBUILD_LAST_N_DAYS=1 /usr/local/sbin/syslog_curate_rotate.sh

この形にしておけば、Iceberg 投入時に ts - INTERVAL 9 HOURS のような補正は不要になります。

13-1. 前日分だけ挿入

syslog:

sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
SET spark.sql.hive.convertMetastoreParquet=false;
SET spark.sql.session.timeZone=Asia/Tokyo;

DELETE FROM hive_prod.logs.syslog_iceberg
WHERE dt = date_sub(current_date(), 1);

INSERT INTO hive_prod.logs.syslog_iceberg
WITH src AS (
  SELECT
    host,
    CAST(ts_text AS TIMESTAMP_NTZ) AS ts_fixed,
    severity,
    program,
    msg,
    CAST(dt AS date) AS dt_fixed
  FROM logs.syslog_curated
  WHERE dt = date_sub(current_date(), 1)
    AND ts_text IS NOT NULL
)
SELECT
  host,
  ts_fixed AS ts,
  severity,
  program,
  msg,
  dt_fixed AS dt,
  HOUR(ts_fixed) AS hr
FROM src;
EOF

authlog:

sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
SET spark.sql.hive.convertMetastoreParquet=false;
SET spark.sql.session.timeZone=Asia/Tokyo;

DELETE FROM hive_prod.logs.authlog_iceberg
WHERE dt = date_sub(current_date(), 1);

INSERT INTO hive_prod.logs.authlog_iceberg
WITH src AS (
  SELECT
    host,
    CAST(ts_text AS TIMESTAMP_NTZ) AS ts_fixed,
    severity,
    program,
    msg,
    CAST(dt AS date) AS dt_fixed
  FROM authlogs.authlog_curated
  WHERE dt = date_sub(current_date(), 1)
    AND ts_text IS NOT NULL
)
SELECT
  host,
  ts_fixed AS ts,
  severity,
  program,
  msg,
  dt_fixed AS dt,
  HOUR(ts_fixed) AS hr
FROM src;
EOF

14. 件数比較

件数が一致すること。

syslog:

sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
SELECT 'curated' AS src, count(*) AS cnt FROM logs.syslog_curated WHERE dt = date_sub(current_date(), 1)
UNION ALL
SELECT 'iceberg' AS src, count(*) AS cnt FROM hive_prod.logs.syslog_iceberg WHERE dt = date_sub(current_date(), 1);
EOF

authlog:

sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
SELECT 'curated' AS src, count(*) AS cnt FROM authlogs.authlog_curated WHERE dt = date_sub(current_date(), 1)
UNION ALL
SELECT 'iceberg' AS src, count(*) AS cnt FROM hive_prod.logs.authlog_iceberg WHERE dt = date_sub(current_date(), 1);
EOF

ope1設定

既存のope1(基盤操作ホスト)へSpark/Icebergメンテ操作できるようにします。

以下の記事を参考にope1を構築していること。

全体像

ope1
├ Spark client(軽量)
├ cron / systemd.timer
└ spark-sql 実行

spark://iceberg1:7077

iceberg1/2/3(実行)

HDFS + Hive Metastore

1. ope1 に Spark client を入れる

cd /tmp
curl -LO https://archive.apache.org/dist/spark/spark-3.5.8/spark-3.5.8-bin-hadoop3.tgz

sudo mkdir -p /opt/spark
sudo tar -xzf spark-3.5.8-bin-hadoop3.tgz -C /opt/spark
sudo ln -sfn /opt/spark/spark-3.5.8-bin-hadoop3 /opt/spark/current

2. Iceberg runtime jar

sudo mkdir -p /opt/spark/extra-jars
cd /opt/spark/extra-jars

sudo curl -LO https://repo1.maven.org/maven2/org/apache/iceberg/iceberg-spark-runtime-3.5_2.12/1.10.1/iceberg-spark-runtime-3.5_2.12-1.10.1.jar

3. Spark 環境設定

sudo mkdir -p /etc/spark/conf
sudo tee /etc/spark/conf/spark-env.sh > /dev/null <<'EOF'
export JAVA_HOME=$(dirname $(dirname $(readlink -f $(which java))))
export SPARK_HOME=/opt/spark/current
export SPARK_CONF_DIR=/etc/spark/conf
export HADOOP_CONF_DIR=/etc/hadoop/conf
export HIVE_CONF_DIR=/etc/hive/conf
EOF

spark-defaults.conf作成

sudo tee /etc/spark/conf/spark-defaults.conf > /dev/null <<'EOF'
spark.master=spark://iceberg1:7077

spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions

spark.sql.catalog.hive_prod=org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.hive_prod.type=hive
spark.sql.catalog.hive_prod.uri=thrift://hive1:9083,thrift://hive2:9083
spark.sql.catalog.hive_prod.warehouse=hdfs://cluster1/warehouse/iceberg

spark.driver.extraClassPath=/opt/spark/extra-jars/iceberg-spark-runtime-3.5_2.12-1.10.1.jar
spark.executor.extraClassPath=/opt/spark/extra-jars/iceberg-spark-runtime-3.5_2.12-1.10.1.jar

spark.sql.sources.partitionOverwriteMode=dynamic
spark.sql.session.timeZone=Asia/Tokyo
EOF
sudo tee /usr/local/bin/spark-sql-iceberg > /dev/null <<'EOF'
#!/usr/bin/env bash
set -euo pipefail

export JAVA_HOME=/usr/lib/jvm/java-17-openjdk-amd64
export PATH="${JAVA_HOME}/bin:/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin"
export SPARK_HOME=/opt/spark/current
export SPARK_CONF_DIR=/etc/spark/conf
export HADOOP_CONF_DIR=/etc/hadoop/conf
export HIVE_CONF_DIR=/etc/hive/conf

exec /opt/spark/current/bin/spark-sql \
  --conf spark.driver.memory=1g \
  --conf spark.executor.memory=512m \
  --conf spark.executor.cores=1 \
  --conf spark.cores.max=1 \
  --conf spark.sql.shuffle.partitions=1 \
  --conf spark.sql.parquet.enableDictionary=false \
  "$@"
EOF

sudo chmod 755 /usr/local/bin/spark-sql-iceberg
sudo useradd -r -m -d /var/lib/spark -s /bin/bash spark || true

sudo mkdir -p /var/lib/spark
sudo chown -R spark:spark /var/lib/spark

4. 動作確認

sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
SHOW DATABASES;
EOF

sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
SHOW TABLES IN hive_prod.logs;
EOF

Icebergからテーブル見えればOK

5. 定期ジョブ用ディレクトリ

sudo mkdir -p /opt/iceberg/bin
sudo chown -R $(whoami):$(whoami) /opt/iceberg

6. 日次投入スクリプト

従来のcureted取り込みが終わったあとに実施する。

仕様

  • 引数でyyyy-mm-ddを指定すれば特定日の取り込み直しを実施
  • 指定しないと前日のデータ取り込みを実施
  • syslog/authlog取り込みを実施
  • 最後にcurated/icebergテーブル件数を比較し、OK/NG判断実施
tee /opt/iceberg/bin/load_curated_to_iceberg.sh <<'EOF'
#!/usr/bin/env bash
set -euo pipefail

DT="${1:-$(date -d 'yesterday' +%F)}"
SPARK_SQL="sudo -u spark /usr/local/bin/spark-sql-iceberg"

log() {
  echo "[INFO] $(date '+%F %T') $*"
}

err() {
  echo "[ERROR] $(date '+%F %T') $*" >&2
}

extract_last_integer() {
  awk '
    {
      gsub(/^[[:space:]]+|[[:space:]]+$/, "", $0)
      if ($0 ~ /^[0-9]+$/) val=$0
    }
    END {
      if (val == "") exit 1
      print val
    }
  '
}

run_spark_count() {
  local sql="$1"
  local out rc

  set +e
  out=$(${SPARK_SQL} 2>&1 <<SQL
${sql}
SQL
)
  rc=$?
  set -e

  if [ "${rc}" -ne 0 ]; then
    printf '%s\n' "${out}" >&2
    err "spark-sql failed"
    return 1
  fi

  printf '%s\n' "${out}" | extract_last_integer
}

log "reload start dt=${DT}"

${SPARK_SQL} <<SQL

SET spark.sql.hive.convertMetastoreParquet=false;
SET spark.sql.session.timeZone=Asia/Tokyo;

REFRESH TABLE logs.syslog_curated;

DELETE FROM hive_prod.logs.syslog_iceberg
WHERE dt = DATE '${DT}';

INSERT INTO hive_prod.logs.syslog_iceberg
WITH src AS (
  SELECT
    host,
    CAST(ts_text AS TIMESTAMP_NTZ) AS ts_fixed,
    severity,
    program,
    msg,
    CAST(dt AS date) AS dt_fixed
  FROM logs.syslog_curated
  WHERE dt = '${DT}'
    AND ts_text IS NOT NULL
)
SELECT
  host,
  ts_fixed AS ts,
  severity,
  program,
  msg,
  dt_fixed AS dt,
  HOUR(ts_fixed) AS hr
FROM src;
SQL

log "syslog reload done dt=${DT}"

${SPARK_SQL} <<SQL

SET spark.sql.hive.convertMetastoreParquet=false;
SET spark.sql.session.timeZone=Asia/Tokyo;

REFRESH TABLE authlogs.authlog_curated;

DELETE FROM hive_prod.logs.authlog_iceberg
WHERE dt = DATE '${DT}';

INSERT INTO hive_prod.logs.authlog_iceberg
WITH src AS (
  SELECT
    host,
    CAST(ts_text AS TIMESTAMP_NTZ) AS ts_fixed,
    severity,
    program,
    msg,
    CAST(dt AS date) AS dt_fixed
  FROM authlogs.authlog_curated
  WHERE dt = '${DT}'
    AND ts_text IS NOT NULL
)
SELECT
  host,
  ts_fixed AS ts,
  severity,
  program,
  msg,
  dt_fixed AS dt,
  HOUR(ts_fixed) AS hr
FROM src;
SQL

log "authlog reload done dt=${DT}"
log "start count check dt=${DT}"

# syslog: Hive curated vs Iceberg を Spark で比較
SYSLOG_HIVE_COUNT="$(run_spark_count "
SELECT COUNT(*)
FROM logs.syslog_curated
WHERE dt = '${DT}';
")"

SYSLOG_ICEBERG_COUNT="$(run_spark_count "
SELECT COUNT(*)
FROM hive_prod.logs.syslog_iceberg
WHERE dt = DATE '${DT}';
")"

SYSLOG_DIFF=$((SYSLOG_ICEBERG_COUNT - SYSLOG_HIVE_COUNT))

log "syslog hive_count=${SYSLOG_HIVE_COUNT}"
log "syslog iceberg_count=${SYSLOG_ICEBERG_COUNT}"
log "syslog diff=${SYSLOG_DIFF}"

# authlog: Hive curated vs Iceberg を Spark で比較
AUTHLOG_HIVE_COUNT="$(run_spark_count "
SELECT COUNT(*)
FROM authlogs.authlog_curated
WHERE dt = '${DT}';
")"

AUTHLOG_ICEBERG_COUNT="$(run_spark_count "
SELECT COUNT(*)
FROM hive_prod.logs.authlog_iceberg
WHERE dt = DATE '${DT}';
")"

AUTHLOG_DIFF=$((AUTHLOG_ICEBERG_COUNT - AUTHLOG_HIVE_COUNT))

log "authlog hive_count=${AUTHLOG_HIVE_COUNT}"
log "authlog iceberg_count=${AUTHLOG_ICEBERG_COUNT}"
log "authlog diff=${AUTHLOG_DIFF}"

FAILED=0

if [ "${SYSLOG_HIVE_COUNT}" != "${SYSLOG_ICEBERG_COUNT}" ]; then
  err "syslog count mismatch dt=${DT} hive=${SYSLOG_HIVE_COUNT} iceberg=${SYSLOG_ICEBERG_COUNT}"
  FAILED=1
else
  log "OK syslog count matched dt=${DT}"
fi

if [ "${AUTHLOG_HIVE_COUNT}" != "${AUTHLOG_ICEBERG_COUNT}" ]; then
  err "authlog count mismatch dt=${DT} hive=${AUTHLOG_HIVE_COUNT} iceberg=${AUTHLOG_ICEBERG_COUNT}"
  FAILED=1
else
  log "OK authlog count matched dt=${DT}"
fi

if [ "${FAILED}" -ne 0 ]; then
  err "reload finished with mismatch dt=${DT}"
  exit 1
fi

log "reload done dt=${DT}"
EOF

chmod 755 /opt/iceberg/bin/load_curated_to_iceberg.sh

7. compaction スクリプト

定期ジョブ完了後に実施する。
小ファイルに分散されたデータファイルをまとめる動作になります。

仕様

  • Iceberg テーブルのデータファイルを整理・圧縮を実施する
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"
)

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

chmod 755 /opt/iceberg/bin/compact_iceberg.sh

8. snapshot cleanup

毎日実施すること。

仕様

  • 引数でyyyy-mm-ddを指定すれば特定日より前のデータの削除を実施
  • 指定しないと14日より前のデータ削除を実施
  • スナップショットは1日より前のを削除する。
tee /opt/iceberg/bin/expire_iceberg.sh <<'EOF'
#!/usr/bin/env bash
set -euo pipefail

DT="${1:-$(date -d '14 days ago' '+%F')}"
SNAP="${SNAP:-$(date -d '1 days ago' '+%F 00:00:00')}"

SPARK_SQL="${SPARK_SQL:-sudo -u spark /usr/local/bin/spark-sql-iceberg}"
CATALOG="${CATALOG:-hive_prod}"
DB="${DB:-logs}"

TABLES=(
  "syslog_iceberg"
  "authlog_iceberg"
)

run_sql() {
  local sql="$1"

  bash -lc "${SPARK_SQL} <<SQL
${sql}
SQL
"
}

echo "[INFO] daily cleanup start dt=${DT} snap=${SNAP}"

for tbl in "${TABLES[@]}"; do
  echo "[INFO] DELETE ${CATALOG}.${DB}.${tbl}"

  run_sql "
DELETE FROM ${CATALOG}.${DB}.${tbl}
WHERE dt < DATE '${DT}';
"

  echo "[INFO] expire_snapshots ${CATALOG}.${DB}.${tbl}"

  run_sql "
CALL ${CATALOG}.system.expire_snapshots(
  table => '${DB}.${tbl}',
  older_than => TIMESTAMP '${SNAP}',
  clean_expired_metadata => true
);
"

done

echo "[INFO] daily cleanup done dt=${DT} snap=${SNAP}"
EOF

chmod 755 /opt/iceberg/bin/expire_iceberg.sh

9. orphan files cleanup

毎日実施すること。

これを行わないとローテーションで外したデータファイル、削除したスナップショットの実ファイルが削除されないため、実施してください。

仕様

  • 1日より前の論理削除データを物理削除する。
  • 初回使用前に手動でDRY_RUNを使用して正常に動くか確認する。
tee /opt/iceberg/bin/remove_orphan_iceberg.sh <<'EOF'
#!/usr/bin/env bash
set -euo pipefail

# true : 削除候補の確認のみ
# false: 実際に削除
DRY_RUN="${DRY_RUN:-true}"

# 何日前より古い orphan ファイルを対象にするか
ORPHAN_DAYS="${ORPHAN_DAYS:-1}"
ORPHAN_TS="$(date -d "${ORPHAN_DAYS} days ago" '+%F %T')"

SPARK_SQL="${SPARK_SQL:-sudo -u spark /usr/local/bin/spark-sql-iceberg}"
CATALOG="${CATALOG:-hive_prod}"
DB="${DB:-logs}"

TABLES=(
  "syslog_iceberg"
  "authlog_iceberg"
)

run_sql() {
  local sql="$1"

  bash -lc "${SPARK_SQL} <<SQL
${sql}
SQL
"
}

case "${DRY_RUN}" in
  true|false)
    ;;
  *)
    echo "[ERROR] DRY_RUN must be true or false: ${DRY_RUN}" >&2
    exit 1
    ;;
esac

if ! [[ "${ORPHAN_DAYS}" =~ ^[0-9]+$ ]]; then
  echo "[ERROR] ORPHAN_DAYS must be a non-negative integer: ${ORPHAN_DAYS}" >&2
  exit 1
fi

echo "[INFO] orphan cleanup start dry_run=${DRY_RUN} orphan_ts=${ORPHAN_TS}"

for tbl in "${TABLES[@]}"; do
  FULL_TABLE="${CATALOG}.${DB}.${tbl}"

  if [[ "${DRY_RUN}" == "true" ]]; then
    echo "[INFO] orphan dry-run start: ${FULL_TABLE}"

    run_sql "
CALL ${CATALOG}.system.remove_orphan_files(
  table => '${DB}.${tbl}',
  older_than => TIMESTAMP '${ORPHAN_TS}',
  dry_run => true
);
"

    echo "[INFO] orphan dry-run done : ${FULL_TABLE}"
  else
    echo "[INFO] orphan delete start: ${FULL_TABLE}"

    run_sql "
CALL ${CATALOG}.system.remove_orphan_files(
  table => '${DB}.${tbl}',
  older_than => TIMESTAMP '${ORPHAN_TS}'
);
"

    echo "[INFO] orphan delete done : ${FULL_TABLE}"
  fi
done

echo "[INFO] orphan cleanup done dry_run=${DRY_RUN} orphan_ts=${ORPHAN_TS}"
EOF

chmod 755 /opt/iceberg/bin/remove_orphan_iceberg.sh

10. 動作テスト

/opt/iceberg/bin/load_curated_to_iceberg.sh
/opt/iceberg/bin/compact_iceberg.sh
/opt/iceberg/bin/expire_iceberg.sh
DRY_RUN=true ORPHAN_DAYS=7 /opt/iceberg/bin/remove_orphan_iceberg.sh 

11. cron 設定

crontab -e
# 日次ロード(前日分)(毎日5:00)
0 5 * * * /opt/iceberg/bin/load_curated_to_iceberg.sh >> /var/log/iceberg_load.log 2>&1

# compaction(毎日5:10)
10 5 * * * /opt/iceberg/bin/compact_iceberg.sh >> /var/log/iceberg_compact.log 2>&1

# snapshot cleanup(毎日5:20)
20 5 * * 0 /opt/iceberg/bin/expire_iceberg.sh >> /var/log/iceberg_expire.log 2>&1

# orphan files cleanup(毎日5:30)
30 5 * * 0 DRY_RUN=false ORPHAN_DAYS=7 /opt/iceberg/bin/remove_orphan_iceberg.sh >> /var/log/remove_orphan_iceberg.log 2>&1

Zeppelin設定

既存のzeppelin1(解析ホスト)へSpark/Iceberg解析操作できるようにします。

以下の記事を参考にzeppelin1を構築していること。

1. Spark client 導入

cd /tmp
curl -LO https://archive.apache.org/dist/spark/spark-3.5.8/spark-3.5.8-bin-hadoop3.tgz

sudo mkdir -p /opt/spark
sudo tar -xzf spark-3.5.8-bin-hadoop3.tgz -C /opt/spark
sudo ln -sfn /opt/spark/spark-3.5.8-bin-hadoop3 /opt/spark/current

2. Icebergruntimejar

sudo mkdir -p /opt/spark/extra-jars
cd /opt/spark/extra-jars

sudo curl -LO https://repo1.maven.org/maven2/org/apache/iceberg/iceberg-spark-runtime-3.5_2.12/1.10.1/iceberg-spark-runtime-3.5_2.12-1.10.1.jar

3. Hadoop / Hive 設定コピー

既存クラスタからコピー

sudo mkdir -p /etc/hadoop/conf
sudo mkdir -p /etc/hive/conf

配置:

/etc/hadoop/conf/core-site.xml
/etc/hadoop/conf/hdfs-site.xml
/etc/hive/conf/hive-site.xml

権限:

sudo chmod 644 /etc/hadoop/conf/*.xml /etc/hive/conf/*.xml

4. Spark 環境設定

sudo mkdir -p /etc/spark/conf
sudo tee /etc/spark/conf/spark-env.sh > /dev/null <<'EOF'
export JAVA_HOME=$(dirname $(dirname $(readlink -f $(which java))))
export SPARK_HOME=/opt/spark/current
export SPARK_CONF_DIR=/etc/spark/conf
export HADOOP_CONF_DIR=/etc/hadoop/conf
export HIVE_CONF_DIR=/etc/hive/conf
EOF
sudo tee /etc/spark/conf/spark-defaults.conf > /dev/null <<'EOF'
spark.master=spark://iceberg1:7077

spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions

spark.sql.catalog.hive_prod=org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.hive_prod.type=hive
spark.sql.catalog.hive_prod.uri=thrift://hive1:9083,thrift://hive2:9083
spark.sql.catalog.hive_prod.warehouse=hdfs://cluster1/warehouse/iceberg

spark.driver.extraClassPath=/opt/spark/extra-jars/iceberg-spark-runtime-3.5_2.12-1.10.1.jar
spark.executor.extraClassPath=/opt/spark/extra-jars/iceberg-spark-runtime-3.5_2.12-1.10.1.jar

spark.sql.sources.partitionOverwriteMode=dynamic
spark.sql.session.timeZone=Asia/Tokyo
EOF
sudo tee /usr/local/bin/spark-sql-iceberg > /dev/null <<'EOF'
#!/usr/bin/env bash
set -euo pipefail

export JAVA_HOME=/usr/lib/jvm/java-11-openjdk-amd64
export PATH="${JAVA_HOME}/bin:/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin"
export SPARK_HOME=/opt/spark/current
export SPARK_CONF_DIR=/etc/spark/conf
export HADOOP_CONF_DIR=/etc/hadoop/conf
export HIVE_CONF_DIR=/etc/hive/conf

exec /opt/spark/current/bin/spark-sql \
  --conf spark.driver.memory=1g \
  --conf spark.executor.memory=512m \
  --conf spark.executor.cores=1 \
  --conf spark.cores.max=1 \
  --conf spark.sql.shuffle.partitions=1 \
  --conf spark.sql.parquet.enableDictionary=false \
  "$@"
EOF

sudo chmod 755 /usr/local/bin/spark-sql-iceberg
sudo useradd -r -m -d /var/lib/spark -s /bin/bash spark || true

sudo mkdir -p /var/lib/spark
sudo chown -R spark:spark /var/lib/spark

5. 単体動作確認

sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
SHOW DATABASES;
EOF

sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
SHOW TABLES IN hive_prod.logs;
EOF

6. Zeppelin Interpreter 設定

フォルダ権限設定

sudo mkdir -p /var/lib/spark/work
sudo chown -R spark:spark /var/lib/spark/work

sudo usermod -aG spark zeppelin
sudo chgrp spark /var/lib/spark/work
sudo chmod 2775 /var/lib/spark/work

Zeppelin UI:

Interpreter → spark → Edit

画面設定を見比べ以下の設定を行う

既存からの書き換え

SPARK_HOME=/opt/spark/current
spark.master=spark://iceberg1:7077
spark.submit.deployMode=client
spark.app.name=Zeppelin-Iceberg
spark.driver.memory=1g
spark.executor.memory=512m
zeppelin.spark.enableSupportedVersionCheck=false
spark.jars=/opt/spark/extra-jars/iceberg-spark-runtime-3.5_2.12-1.10.1.jar

以下のパラメータについて、削除して入れ直す。

spark.executor.instances=1

以下のパラメータを追加する。

HADOOP_CONF_DIR=/etc/hadoop/conf
HIVE_CONF_DIR=/etc/hive/conf
SPARK_CONF_DIR=/etc/spark/conf
spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions
spark.sql.catalog.hive_prod=org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.hive_prod.type=hive
spark.sql.catalog.hive_prod.uri=thrift://hive1:9083,thrift://hive2:9083
spark.sql.catalog.hive_prod.warehouse=hdfs://cluster1/warehouse/iceberg
spark.sql.sources.partitionOverwriteMode=dynamic
spark.sql.session.timeZone=Asia/Tokyo
spark.cores.max=1
spark.sql.shuffle.partitions=1
spark.sql.parquet.enableDictionary=false
spark.sql.catalogImplementation=hive
hive.metastore.uris=thrift://hive1:9083,thrift://hive2:9083

Save
restart

zeppelinの再起動も実施する。

sudo systemctl restart zeppelin

7. Zeppelin で確認

%spark.sql

SHOW TABLES IN hive_prod.logs;

SELECT count(*) FROM hive_prod.logs.syslog_iceberg;

8. zeppelinでデータ編集する際の設定

zeppelinからのデータ編集する際は、
以下のホストで操作をしてください。

ope1:

sudo -u hadoop hdfs dfs -setfacl -R -m user:zeppelin:rwx /warehouse/iceberg/logs/syslog_iceberg
sudo -u hadoop hdfs dfs -setfacl -R -m default:user:zeppelin:rwx /warehouse/iceberg/logs/syslog_iceberg
sudo -u hadoop hdfs dfs -setfacl -R -m user:zeppelin:rwx /warehouse/iceberg/logs/authlog_iceberg
sudo -u hadoop hdfs dfs -setfacl -R -m default:user:zeppelin:rwx /warehouse/iceberg/logs/authlog_iceberg

sudo -u hadoop hdfs dfs -setfacl -R -m user:spark:rwx /warehouse/iceberg/logs/syslog_iceberg
sudo -u hadoop hdfs dfs -setfacl -R -m default:user:spark:rwx /warehouse/iceberg/logs/syslog_iceberg
sudo -u hadoop hdfs dfs -setfacl -R -m user:spark:rwx /warehouse/iceberg/logs/authlog_iceberg
sudo -u hadoop hdfs dfs -setfacl -R -m default:user:spark:rwx /warehouse/iceberg/logs/authlog_iceberg

解析SQL

以下のSQLで解析を行えます。

0. curated vs iceberg 差分(差分が出る場合要調査)

%spark.sql
SELECT 'curated' AS src, count(*) AS cnt FROM logs.syslog_curated
UNION ALL
SELECT 'iceberg' AS src, count(*) AS cnt FROM hive_prod.logs.syslog_iceberg;

SELECT 'curated' AS src, count(*) AS cnt FROM authlogs.authlog_curated
UNION ALL
SELECT 'iceberg' AS src, count(*) AS cnt FROM hive_prod.logs.authlog_iceberg;

1. 基本:件数推移(日次)

%spark.sql

SELECT
  dt,
  COUNT(*) AS cnt
FROM hive_prod.logs.syslog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY dt
ORDER BY dt;

SELECT
  dt,
  COUNT(*) AS cnt
FROM hive_prod.logs.authlog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY dt
ORDER BY dt;

2. ホスト別ログ量

%spark.sql
SELECT
  host,
  COUNT(*) AS cnt
FROM hive_prod.logs.syslog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY host
ORDER BY cnt DESC;

SELECT
  host,
  COUNT(*) AS cnt
FROM hive_prod.logs.authlog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY host
ORDER BY cnt DESC;

3. severity 分布

%spark.sql
SELECT
  severity,
  COUNT(*) AS cnt
FROM hive_prod.logs.syslog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY severity
ORDER BY severity;

SELECT
  severity,
  COUNT(*) AS cnt
FROM hive_prod.logs.authlog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY severity
ORDER BY severity;

4. program別

%spark.sql
SELECT
  program,
  COUNT(*) AS cnt
FROM hive_prod.logs.syslog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY program
ORDER BY cnt DESC;

SELECT
  program,
  COUNT(*) AS cnt
FROM hive_prod.logs.authlog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY program
ORDER BY cnt DESC;

5. エラー系ログ抽出

%spark.sql
SELECT host,ts,severity,program,msg
FROM hive_prod.logs.syslog_iceberg
WHERE dt >= date_sub(current_date(), 14)
  AND severity <= 3
ORDER BY ts DESC
LIMIT 100;

SELECT host,ts,severity,program,msg
FROM hive_prod.logs.authlog_iceberg
WHERE dt >= date_sub(current_date(), 14)
  AND severity <= 3
ORDER BY ts DESC
LIMIT 100;

6. keyword検索(障害調査)

%spark.sql
SELECT host,ts,severity,program,msg
FROM hive_prod.logs.syslog_iceberg
WHERE dt >= date_sub(current_date(), 14)
  AND msg LIKE '%timeout%'
ORDER BY ts DESC
LIMIT 100;

SELECT host,ts,severity,program,msg
FROM hive_prod.logs.authlog_iceberg
WHERE dt >= date_sub(current_date(), 14)
  AND msg LIKE '%timeout%'
ORDER BY ts DESC
LIMIT 100;

7. 成功ログイン抽出(authlog)

%spark.sql
SSELECT host,ts,severity,program,msg
FROM hive_prod.logs.authlog_iceberg
WHERE dt = date_sub(current_date(), 1)
  AND msg LIKE '%Accepted publickey%'
ORDER BY ts DESC
LIMIT 100;

8. 成功ログインuser抽出(authlog)

%spark.sql
SELECT
  regexp_extract(msg, 'for ([^ ]+)', 1) AS user,
  COUNT(*) AS cnt
FROM hive_prod.logs.authlog_iceberg
WHERE dt = date_sub(current_date(), 2)
  AND msg LIKE '%Accepted publickey%'
GROUP BY regexp_extract(msg, 'for ([^ ]+)', 1)
ORDER BY cnt DESC;

9. syslog × authlog 相関(時間帯)

%spark.sql
SELECT
  s.x_time,
  s.syslog_cnt,
  a.auth_cnt
FROM (
  SELECT
    concat(CAST(dt AS string), ' ', lpad(CAST(hr AS string), 2, '0')) AS x_time,
    COUNT(*) AS syslog_cnt
  FROM hive_prod.logs.syslog_iceberg
  GROUP BY dt, hr
) s
JOIN (
  SELECT
    concat(CAST(dt AS string), ' ', lpad(CAST(hr AS string), 2, '0')) AS x_time,
    COUNT(*) AS auth_cnt
  FROM hive_prod.logs.authlog_iceberg
  GROUP BY dt, hr
) a
ON s.x_time = a.x_time
ORDER BY s.x_time;

10. snapshot参照

数字の羅列部分はsnapshot_idを指定する。
これを指定することによりこの時点でのデータを参照可能

%spark.sql

SELECT snapshot_id, committed_at, operation
FROM hive_prod.logs.syslog_iceberg.snapshots
ORDER BY committed_at DESC;

SELECT *
FROM hive_prod.logs.syslog_iceberg
VERSION AS OF 4046019217212690338
ORDER BY ts DESC
LIMIT 10;

SELECT snapshot_id, committed_at, operation
FROM hive_prod.logs.authlog_iceberg.snapshots
ORDER BY committed_at DESC;

SELECT *
FROM hive_prod.logs.authlog_iceberg
VERSION AS OF 2646869706553548116
ORDER BY ts DESC
LIMIT 10;

11. 最近の異常急増検知

平均の2倍以上のログが出たときに検知

%spark.sql
SELECT
  dt,
  COUNT(*) AS cnt
FROM hive_prod.logs.syslog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY dt
HAVING cnt > (
  SELECT AVG(cnt) * 2 FROM (
    SELECT COUNT(*) AS cnt
    FROM hive_prod.logs.syslog_iceberg
    GROUP BY dt
  )
)
ORDER BY dt;

SELECT
  dt,
  COUNT(*) AS cnt
FROM hive_prod.logs.authlog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY dt
HAVING cnt > (
  SELECT AVG(cnt) * 2 FROM (
    SELECT COUNT(*) AS cnt
    FROM hive_prod.logs.authlog_iceberg
    GROUP BY dt
  )
)
ORDER BY dt;

12. 日単位収集(syslog)

SELECT
  host,
  dt AS x_day,
  COUNT(*) AS cnt
FROM hive_prod.logs.syslog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY
  host,
  dt
ORDER BY
  x_day,
  host;

13. 日単位収集(authlog)

%spark.sql
SELECT
  host,
  dt AS x_day,
  COUNT(*) AS cnt
FROM hive_prod.logs.authlog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY
  host,
  dt
ORDER BY
  x_day,
  host;

14. 時間単位収集(syslog)

%spark.sql
SELECT
  host,
  concat(
    CAST(dt AS string),
    ' ',
    lpad(CAST(hr AS string), 2, '0')
  ) AS x_time,
  COUNT(*) AS cnt
FROM hive_prod.logs.syslog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY dt, hr, host
ORDER BY dt, hr, host;

15. 時間単位収集(authlog)

SELECT
  host,
  concat(
    CAST(dt AS string),
    ' ',
    lpad(CAST(hr AS string), 2, '0')
  ) AS x_time,
  COUNT(*) AS cnt
FROM hive_prod.logs.authlog_iceberg
WHERE dt >= date_sub(current_date(), 14)
GROUP BY dt, hr, host
ORDER BY dt, hr, host;

性能比較

Hive curated vs Iceberg を同条件で比較
Icebergを使用することによりどのくらい改善するか確認する。

  • ope1から実行
  • CSVで結果保存

0. 事前前提(確認)

hive -e "SELECT 1;"
sudo -u spark /usr/local/bin/spark-sql-iceberg <<'EOF'
SHOW DATABASES;
EOF

1. ディレクトリ作成

sudo mkdir -p /opt/bench
sudo mkdir -p /opt/bench/result
sudo chown -R $(whoami):$(whoami) /opt/bench

2. ベンチスクリプト作成

cat <<'EOF' > /opt/bench/run_benchmark.sh
#!/usr/bin/env bash
set -euo pipefail

DT="${1:-$(date -d 'yesterday' +%F)}"
OUT_DIR="/opt/bench/result"
mkdir -p "${OUT_DIR}"

timestamp() {
  date '+%Y-%m-%d %H:%M:%S'
}

log() {
  echo "[$(timestamp)] $*"
}

run_hive() {
  local name="$1"
  local sql="$2"

  log "HIVE ${name} start"

  local start=$(date +%s)

  hive -e "${sql}" > "${OUT_DIR}/hive_${name}.out" 2>&1

  local end=$(date +%s)
  local elapsed=$((end - start))

  log "HIVE ${name} done ${elapsed}s"

  echo "HIVE,${name},${DT},${elapsed}" >> "${OUT_DIR}/summary.csv"
}

run_spark() {
  local name="$1"
  local sql="$2"

  log "ICEBERG ${name} start"

  local start=$(date +%s)

  sudo -u spark /usr/local/bin/spark-sql-iceberg -e "${sql}" > "${OUT_DIR}/iceberg_${name}.out" 2>&1

  local end=$(date +%s)
  local elapsed=$((end - start))

  log "ICEBERG ${name} done ${elapsed}s"

  echo "ICEBERG,${name},${DT},${elapsed}" >> "${OUT_DIR}/summary.csv"
}

echo "engine,query,dt,sec" > "${OUT_DIR}/summary.csv"

# ========================
# SQL定義
# ========================

HIVE_SET_COMMON=$(cat <<'SQL'
SET mapreduce.map.memory.mb=512;
SET mapreduce.reduce.memory.mb=512;
set yarn.app.mapreduce.am.resource.mb=512;
SQL
)

HIVE_SQL_COUNT="${HIVE_SET_COMMON}
SELECT COUNT(*) FROM logs.syslog_curated WHERE dt='${DT}';"

ICEBERG_SQL_COUNT="SELECT COUNT(*) FROM hive_prod.logs.syslog_iceberg WHERE dt=DATE '${DT}'"

HIVE_SQL_GROUP="${HIVE_SET_COMMON}
SELECT host, COUNT(*) FROM logs.syslog_curated WHERE dt='${DT}' GROUP BY host;"

ICEBERG_SQL_GROUP="SELECT host, COUNT(*) FROM hive_prod.logs.syslog_iceberg WHERE dt=DATE '${DT}' GROUP BY host"

HIVE_SQL_TIME="${HIVE_SET_COMMON}
SELECT hour(ts) AS hr, COUNT(*) FROM logs.syslog_curated WHERE dt='${DT}' GROUP BY hour(ts) ORDER BY hr;"

ICEBERG_SQL_TIME="SELECT hr, COUNT(*) FROM hive_prod.logs.syslog_iceberg WHERE dt=DATE '${DT}' GROUP BY hr"

# ========================
# 実行
# ========================

run_hive "count" "${HIVE_SQL_COUNT}"
run_spark "count" "${ICEBERG_SQL_COUNT}"

run_hive "group" "${HIVE_SQL_GROUP}"
run_spark "group" "${ICEBERG_SQL_GROUP}"

run_hive "time" "${HIVE_SQL_TIME}"
run_spark "time" "${ICEBERG_SQL_TIME}"

# ========================
# ファイル数比較
# ========================

log "FILE COUNT CHECK"

HDFS_PATH="/data/kafka/syslog"

if hdfs dfs -test -d "${HDFS_PATH}"; then
  HIVE_FILE_COUNT=$(
    hdfs dfs -ls -R "${HDFS_PATH}" 2>/dev/null \
      | grep '\.json$' \
      | wc -l
  )
else
  HIVE_FILE_COUNT=0
fi

log "HIVE_FILE_COUNT done ${HIVE_FILE_COUNT}"

echo "HIVE_FILE_COUNT,${HIVE_FILE_COUNT}" >> "${OUT_DIR}/summary.csv"

ICEBERG_FILE_COUNT=$(
sudo -u spark /usr/local/bin/spark-sql-iceberg 2>/dev/null <<SQL \
| awk '/^[0-9]+$/ {print $1}' | tail -1
SELECT COUNT(*)
FROM hive_prod.logs.syslog_iceberg.files;
SQL
)

ICEBERG_FILE_COUNT="${ICEBERG_FILE_COUNT:-0}"

log "ICEBERG_FILE_COUNT done ${ICEBERG_FILE_COUNT}"

echo "ICEBERG_FILE_COUNT,${ICEBERG_FILE_COUNT}" >> "${OUT_DIR}/summary.csv"

log "DONE"
EOF

3. 実行権限

chmod 755 /opt/bench/run_benchmark.sh

4. 実行

/opt/bench/run_benchmark.sh

日付指定も可:

/opt/bench/run_benchmark.sh yyyy-mm-dd

5. 結果確認

cat /opt/bench/result/summary.csv

例:

engine,query,dt,sec
HIVE,count,2026-03-21,12
ICEBERG,count,2026-03-21,5
HIVE,group,2026-03-21,18
ICEBERG,group,2026-03-21,7
HIVE,time,2026-03-21,20
ICEBERG,time,2026-03-21,6
HIVE_FILE_COUNT,3500
ICEBERG_FILE_COUNT,120

6. 正確性チェック

DT=$(date -d 'yesterday' +%F)

HIVE_SET_COMMON=$(cat <<'SQL'
SET mapreduce.map.memory.mb=512;
SET mapreduce.reduce.memory.mb=512;
set yarn.app.mapreduce.am.resource.mb=512;
SQL
)

hive -e "${HIVE_SET_COMMON} SELECT COUNT(*) FROM logs.syslog_curated WHERE dt='${DT}';"

sudo -u spark /usr/local/bin/spark-sql-iceberg -e \
"SELECT COUNT(*) FROM hive_prod.logs.syslog_iceberg WHERE dt=DATE '${DT}'"

必ず一致すること

7. compaction比較

# BEFORE
/opt/bench/run_benchmark.sh

# compaction
/opt/iceberg/bin/compact_iceberg.sh

# AFTER
/opt/bench/run_benchmark.sh

8. 見るべきポイント

① 実行時間
ICEBERG < HIVE
② ファイル数
ICEBERG << HIVE
③ compaction後
状況によりICEBERGのファイル数が少なくなること。

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?