はじめに
完成させたログ収集・解析基盤に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上のApache IcebergをAWS S3へ複製してAthenaで分析する構成を考える【設計編】
- 自宅HDFS上のApache IcebergをAmazon S3へ複製しAthenaから分析してみる【構築編】
- AWS S3上のApache IcebergをAthena向けに運用する【運用編】
- AWS S3上のApache IcebergをAthena経由で可視化する【発展編】
ポイントは、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_text を TIMESTAMP_NTZ に変換します。
まず ts_raw と ts_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_curated に ts_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のファイル数が少なくなること。