はじめに
これまでの記事では、自宅 HDFS 上の Apache Iceberg を Source of Truth として残し、AWS 側に S3 + Glue Data Catalog + Athena の分析用コピーを作りました。
- 自宅HDFS上のApache IcebergをAWS S3へ複製してAthenaで分析する構成を考える【設計編】
- 自宅HDFS上のApache IcebergをAmazon S3へ複製しAthenaから分析してみる【構築編】
- AWS S3上のApache IcebergをAthena向けに運用する【運用編】
- AWS S3上のApache IcebergをAthena経由で可視化する【発展編】
今回は逆方向です。
AWS S3 Iceberg
↓
Spark
↓
自宅 HDFS Iceberg
ただし、この記事で扱うのは双方向同期ではありません。
AWS 側の Iceberg コピーを使って、自宅側の特定日データを戻すためのリストア手順です。
基本方針は変えません。
自宅 = メイン環境 / Source of Truth
AWS = 分析用コピー / 復旧時の参照元
AWS 側を正本に昇格させるのではなく、自宅側で特定日の partition 相当データを壊した、消した、再作成したい、というケースで使う想定です。
この記事は、少なくとも構築編まで完了している環境を前提にします。
ゴール
AWS 側 Iceberg テーブルから、自宅側 Iceberg テーブルへ対象日だけ戻します。
Spark Restore Job は自宅側の spark ユーザーで実行します。
AWS 側の glue_prod catalog から S3 Iceberg テーブルを読み、自宅側の hive_prod catalog へ書き戻します。
リストア対象は、これまでのシリーズと同じ2テーブルです。
glue_prod.${GLUE_DATABASE}.syslog_iceberg -> hive_prod.logs.syslog_iceberg
glue_prod.${GLUE_DATABASE}.authlog_iceberg -> hive_prod.logs.authlog_iceberg
対象日は dt で指定します。
DT=2026-08-18
注意点
この手順は自宅側 Iceberg テーブルへ書き込みます。
通常の AWS 向けエクスポートよりも影響が大きいです。
特に以下に注意します。
- 自宅側を Source of Truth とする方針は維持する
- 常時の双方向同期にはしない
- AWS 側に同期済みの範囲だけ戻せる
- 対象日を間違えると自宅側の既存データを削除して入れ直す
- 通常同期や手動再実行と同じ
dtで同時に動かさない -
syslog_icebergとauthlog_icebergの schema を決め打ちしない - timestamp 型、partition spec、table property の差異を確認する
- 実行前後で件数比較する
今回の方式は、対象日の自宅側データを削除してから AWS 側データを挿入します。
自宅側の対象日を DELETE
↓
AWS 側 Iceberg から INSERT
↓
件数比較
これにより、同じ DT を再実行しても件数が増え続けにくくなります。
一方で、対象日指定を間違えると自宅側の正本に影響します。
そのため、スクリプトでは CONFIRM_RESTORE=1 を指定した場合だけ実行するようにします。
本シリーズの運用では、対象テーブルは日次または手動バッチ以外では更新しない前提です。
そのため、復旧時は対象 dt の通常同期や手動再実行が動いていないことを確認してから実行します。
この前提であれば、双方向同期のような競合解決は不要にできます。
前提
構築編で作成した Spark 実行環境を使います。
| 項目 | 値 |
|---|---|
| Spark | 3.5.8 |
| Apache Iceberg | 1.10.1 |
| 自宅側 catalog | hive_prod |
| AWS側 catalog | glue_prod |
| 自宅側 database | logs |
| AWS側 database | ${GLUE_DATABASE} |
| AWS profile | home-aws |
| AWS credentials | /var/lib/spark/.aws/credentials |
Spark から、次の両方の catalog を参照できる必要があります。
hive_prod = 自宅 HDFS Iceberg
glue_prod = AWS S3 Iceberg
実行ユーザーは spark とします。
source /etc/profile.d/iceberg-s3-athena.sh
use_iceberg_aws_sync
確認します。
printf 'AWS_PROFILE=%s\n' "${AWS_PROFILE}"
printf 'AWS_SHARED_CREDENTIALS_FILE=%s\n' "${AWS_SHARED_CREDENTIALS_FILE}"
printf 'AWS_CONFIG_FILE=%s\n' "${AWS_CONFIG_FILE}"
printf 'AWS_REGION=%s\n' "${AWS_REGION}"
printf 'GLUE_DATABASE=%s\n' "${GLUE_DATABASE}"
AWS_PROFILE=home-aws、AWS_SHARED_CREDENTIALS_FILE=/var/lib/spark/.aws/credentials になっていることを確認してから進めます。
事前確認
まず、AWS 側と自宅側の両方が Spark SQL から見えることを確認します。
実行ユーザー: spark
spark-sql-iceberg-aws <<EOF
SHOW TABLES IN glue_prod.logs;
SHOW TABLES IN hive_prod.logs;
EOF
schema を確認します。
spark-sql-iceberg-aws <<EOF
DESCRIBE TABLE glue_prod.logs.syslog_iceberg;
DESCRIBE TABLE hive_prod.logs.syslog_iceberg;
DESCRIBE TABLE glue_prod.logs.authlog_iceberg;
DESCRIBE TABLE hive_prod.logs.authlog_iceberg;
EOF
作成 SQL も確認します。
spark-sql-iceberg-aws <<EOF
SHOW CREATE TABLE glue_prod.logs.syslog_iceberg;
SHOW CREATE TABLE hive_prod.logs.syslog_iceberg;
SHOW CREATE TABLE glue_prod.logs.authlog_iceberg;
SHOW CREATE TABLE hive_prod.logs.authlog_iceberg;
EOF
Spark SQL の対話画面に直接貼る場合は、次のような SQL として実行します。
SHOW CREATE TABLE glue_prod.logs.syslog_iceberg;
SHOW CREATE TABLE hive_prod.logs.syslog_iceberg;
SHOW CREATE TABLE glue_prod.logs.authlog_iceberg;
SHOW CREATE TABLE hive_prod.logs.authlog_iceberg;
確認する点です。
- column name
- column type
- column order
- nullable
- partition spec
- table property
-
tsなど timestamp 系の型
特に INSERT INTO ... SELECT * を使う場合、列順と型の差異に注意します。
差異がある場合は、SELECT * ではなく列名を明示します。
件数確認
戻したい日付の件数を確認します。
例では 2026-08-18 を対象にします。
RESTORE_DT=2026-08-18
spark-sql-iceberg-aws <<EOF
SELECT count(*)
FROM glue_prod.logs.syslog_iceberg
WHERE dt = DATE '${RESTORE_DT}';
SELECT count(*)
FROM hive_prod.logs.syslog_iceberg
WHERE dt = DATE '${RESTORE_DT}';
EOF
authlog_iceberg も確認します。
RESTORE_DT=2026-08-18
spark-sql-iceberg-aws <<EOF
SELECT count(*)
FROM glue_prod.logs.authlog_iceberg
WHERE dt = DATE '${RESTORE_DT}';
SELECT count(*)
FROM hive_prod.logs.authlog_iceberg
WHERE dt = DATE '${RESTORE_DT}';
EOF
AWS 側の件数が 0 の場合は、リストア元がありません。
その場合は、対象日が間違っていないか、AWS への同期が完了していた日付かを確認します。
手動で1テーブルだけ戻す
まずは syslog_iceberg を手動で戻す例です。
RESTORE_DT=2026-08-18
spark-sql-iceberg-aws <<EOF
DELETE FROM hive_prod.logs.syslog_iceberg
WHERE dt = DATE '${RESTORE_DT}';
INSERT INTO hive_prod.logs.syslog_iceberg
SELECT *
FROM glue_prod.logs.syslog_iceberg
WHERE dt = DATE '${RESTORE_DT}';
EOF
件数を確認します。
RESTORE_DT=2026-08-18
spark-sql-iceberg-aws <<EOF
SELECT count(*)
FROM glue_prod.logs.syslog_iceberg
WHERE dt = DATE '${RESTORE_DT}';
SELECT count(*)
FROM hive_prod.logs.syslog_iceberg
WHERE dt = DATE '${RESTORE_DT}';
EOF
authlog_iceberg も同じです。
RESTORE_DT=2026-08-18
spark-sql-iceberg-aws <<EOF
DELETE FROM hive_prod.logs.authlog_iceberg
WHERE dt = DATE '${RESTORE_DT}';
INSERT INTO hive_prod.logs.authlog_iceberg
SELECT *
FROM glue_prod.logs.authlog_iceberg
WHERE dt = DATE '${RESTORE_DT}';
EOF
リストアスクリプト
手動 SQL だけだと、対象日や向きを間違えやすいです。
共通スクリプトにします。
ファイル:
/opt/iceberg/bin/restore_iceberg_from_s3.sh
実行ユーザー: 通常ユーザー
sudo mkdir -p /opt/iceberg/bin
sudo tee /opt/iceberg/bin/restore_iceberg_from_s3.sh > /dev/null <<'EOF'
#!/usr/bin/env bash
set -euo pipefail
SOURCE_CATALOG=${SOURCE_CATALOG:-glue_prod}
DEST_CATALOG=${DEST_CATALOG:-hive_prod}
SOURCE_DB=${SOURCE_DB:-logs}
DEST_DB=${DEST_DB:-logs}
TABLE=${TABLE:?TABLE is required}
DT=${DT:?DT is required. Example: DT=2026-08-18}
SPARK_SQL_BIN=${SPARK_SQL_BIN:-/usr/local/bin/spark-sql-iceberg-aws}
CONFIRM_RESTORE=${CONFIRM_RESTORE:-0}
ALLOW_ZERO_SOURCE=${ALLOW_ZERO_SOURCE:-0}
: "${AWS_PROFILE:?AWS_PROFILE is required}"
: "${AWS_SHARED_CREDENTIALS_FILE:?AWS_SHARED_CREDENTIALS_FILE is required}"
: "${AWS_CONFIG_FILE:?AWS_CONFIG_FILE is required}"
: "${AWS_REGION:?AWS_REGION is required}"
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_count() {
local sql="$1"
local out
out="$("${SPARK_SQL_BIN}" -S -e "${sql}")"
printf '%s\n' "${out}" | extract_last_integer
}
main() {
local source_table="${SOURCE_CATALOG}.${SOURCE_DB}.${TABLE}"
local dest_table="${DEST_CATALOG}.${DEST_DB}.${TABLE}"
local source_count
local dest_before_count
local dest_after_count
log "restore from AWS Iceberg copy"
log "TABLE=${TABLE}"
log "DT=${DT}"
log "SOURCE=${source_table}"
log "DEST=${dest_table}"
if [ "${CONFIRM_RESTORE}" != "1" ]; then
err "CONFIRM_RESTORE=1 is required because this deletes and rewrites data in ${dest_table}"
exit 2
fi
source_count="$(run_count "SELECT count(*) FROM ${source_table} WHERE dt = DATE '${DT}'")"
dest_before_count="$(run_count "SELECT count(*) FROM ${dest_table} WHERE dt = DATE '${DT}'")"
log "source count=${source_count}"
log "dest count before restore=${dest_before_count}"
if [ "${source_count}" = "0" ] && [ "${ALLOW_ZERO_SOURCE}" != "1" ]; then
err "source count is zero. Check DT or AWS sync status. Set ALLOW_ZERO_SOURCE=1 only if this is intentional."
exit 3
fi
"${SPARK_SQL_BIN}" <<SQL
DELETE FROM ${dest_table}
WHERE dt = DATE '${DT}';
INSERT INTO ${dest_table}
SELECT *
FROM ${source_table}
WHERE dt = DATE '${DT}';
SQL
dest_after_count="$(run_count "SELECT count(*) FROM ${dest_table} WHERE dt = DATE '${DT}'")"
log "dest count after restore=${dest_after_count}"
if [ "${source_count}" != "${dest_after_count}" ]; then
err "count mismatch: source=${source_count}, dest_after=${dest_after_count}"
exit 1
fi
log "restore completed successfully"
}
main "$@"
EOF
sudo chmod 755 /opt/iceberg/bin/restore_iceberg_from_s3.sh
sudo chown spark:spark /opt/iceberg/bin/restore_iceberg_from_s3.sh
このスクリプトは、以下を確認します。
source count = glue_prod.<GLUE_DATABASE>.<table> の対象 DT 件数
dest before count = hive_prod.logs.<table> の対象 DT 件数
dest after count = hive_prod.logs.<table> のリストア後件数
source count と dest after count が一致しない場合は失敗させます。
1テーブルをリストアする
実行ユーザー: spark
source /etc/profile.d/iceberg-s3-athena.sh
use_iceberg_aws_sync
syslog_iceberg を戻します。
TABLE=syslog_iceberg \
DT=2026-08-18 \
CONFIRM_RESTORE=1 \
/opt/iceberg/bin/restore_iceberg_from_s3.sh
authlog_iceberg を戻します。
TABLE=authlog_iceberg \
DT=2026-08-18 \
CONFIRM_RESTORE=1 \
/opt/iceberg/bin/restore_iceberg_from_s3.sh
syslog/authlogをまとめて戻す
2テーブルを順番に戻す親スクリプトを作ります。
ファイル:
/opt/iceberg/bin/restore_logs_from_s3.sh
実行ユーザー: 通常ユーザー
sudo mkdir -p /opt/iceberg/bin
sudo tee /opt/iceberg/bin/restore_logs_from_s3.sh > /dev/null <<'EOF'
#!/usr/bin/env bash
set -euo pipefail
DT=${DT:?DT is required. Example: DT=2026-08-18}
RESTORE_BIN=${RESTORE_BIN:-/opt/iceberg/bin/restore_iceberg_from_s3.sh}
SPARK_SQL_BIN=${SPARK_SQL_BIN:-/usr/local/bin/spark-sql-iceberg-aws}
CONFIRM_RESTORE=${CONFIRM_RESTORE:-0}
TABLES=(
syslog_iceberg
authlog_iceberg
)
log() {
echo "[INFO] $(date '+%F %T') $*"
}
if [ "${CONFIRM_RESTORE}" != "1" ]; then
echo "[ERROR] CONFIRM_RESTORE=1 is required" >&2
exit 2
fi
log "restore logs from S3 start: DT=${DT}"
for TABLE in "${TABLES[@]}"; do
log "restore start: TABLE=${TABLE}"
TABLE="${TABLE}" \
DT="${DT}" \
SPARK_SQL_BIN="${SPARK_SQL_BIN}" \
CONFIRM_RESTORE="${CONFIRM_RESTORE}" \
"${RESTORE_BIN}"
log "restore finished: TABLE=${TABLE}"
done
log "restore logs from S3 completed successfully: DT=${DT}"
EOF
sudo chmod 755 /opt/iceberg/bin/restore_logs_from_s3.sh
sudo chown spark:spark /opt/iceberg/bin/restore_logs_from_s3.sh
実行します。
実行ユーザー: spark
DT=2026-08-18 \
CONFIRM_RESTORE=1 \
/opt/iceberg/bin/restore_logs_from_s3.sh
この操作は日次 timer には登録しません。
必要なときだけ手動で実行します。
リストア後の確認
自宅側を確認します。
spark-sql-iceberg-aws <<'EOF'
SELECT
dt,
count(*) AS cnt
FROM hive_prod.logs.syslog_iceberg
WHERE dt = DATE '2026-08-18'
GROUP BY dt;
EOF
AWS 側と比較します。
spark-sql-iceberg-aws <<'EOF'
SELECT
dt,
count(*) AS cnt
FROM glue_prod.logs.syslog_iceberg
WHERE dt = DATE '2026-08-18'
GROUP BY dt;
EOF
ホスト別にも確認します。
spark-sql-iceberg-aws <<'EOF'
SELECT
dt,
host,
count(*) AS cnt
FROM hive_prod.logs.syslog_iceberg
WHERE dt = DATE '2026-08-18'
GROUP BY dt, host
ORDER BY host;
EOF
spark-sql-iceberg-aws <<'EOF'
SELECT
dt,
host,
count(*) AS cnt
FROM glue_prod.logs.syslog_iceberg
WHERE dt = DATE '2026-08-18'
GROUP BY dt, host
ORDER BY host;
EOF
authlog_iceberg も同じ粒度で比較します。
失敗時の考え方
source count が 0
AWS 側に対象日のデータがありません。
確認する点です。
-
DTが正しいか - AWS への日次同期が成功していたか
-
SOURCE_DBが正しいか -
GLUE_DATABASEが想定どおりか - テーブル名が正しいか
通常はこの時点で止めます。
空データで自宅側を上書きしたいケースはほぼないためです。
count mismatch
リストア後の自宅側件数と AWS 側件数が一致していません。
確認する点です。
- Spark job が途中で失敗していないか
- HDFS 側で書き込みエラーが出ていないか
- schema 差異がないか
- 対象日の
dt型が両側で同じか -
INSERT INTO ... SELECT *の列順に問題がないか
schema 差異がある場合は、列名を明示します。
RESTORE_DT=2026-08-18
spark-sql-iceberg-aws <<EOF
INSERT INTO hive_prod.logs.syslog_iceberg (
host,
ts,
severity,
program,
msg,
dt,
hr
)
SELECT
host,
ts,
severity,
program,
msg,
dt,
hr
FROM glue_prod.logs.syslog_iceberg
WHERE dt = DATE '${RESTORE_DT}';
EOF
DELETE 後に INSERT が失敗した
Iceberg は snapshot を持つため、自宅側 table の直前 snapshot へ戻せる可能性があります。
ただし、実際に戻すには自宅側 Iceberg の snapshot 保持状態と利用 engine の対応を確認します。
まず snapshot を確認します。
spark-sql-iceberg-aws <<'EOF'
SELECT *
FROM hive_prod.logs.syslog_iceberg.snapshots
ORDER BY committed_at DESC;
EOF
運用編で整理した snapshot 管理を強く削りすぎていると、戻せる範囲が短くなります。
復旧手順を使う可能性がある場合は、snapshot の保持期間も合わせて設計します。
運用ルール
このリストア手順は便利ですが、普段のデータ流れに組み込むものではありません。
運用ルールとしては、次のようにしておきます。
| 項目 | 方針 |
|---|---|
| 通常同期 | 自宅 HDFS Iceberg から AWS S3 Iceberg への片方向 |
| 逆方向 | 障害時、検証時、手動復旧時だけ |
| 実行単位 | 原則 dt 単位 |
| 実行ユーザー | spark |
| 実行条件 |
CONFIRM_RESTORE=1 を明示 |
| 自動化 | systemd timer には登録しない |
| 同時実行 | 同じ dt の通常同期、手動再実行、リストアを並列実行しない |
| 確認 | リストア前後の件数比較を必須にする |
双方向同期として常時動かす場合は、競合解決が必要になります。
例えば、自宅側と AWS 側の同じ dt に別々の更新が入った場合、どちらを正とするかを決めなければなりません。
今回のログ基盤では、そこまで複雑にせず、自宅側を正本、AWS 側を分析用コピー兼リストア元として扱います。
対象テーブルを日次または手動バッチでだけ更新する運用であれば、復旧作業の時間帯を分けることで競合リスクを低く抑えられます。
まとめ
AWS へ複製した Iceberg テーブルは、条件付きで自宅 HDFS Iceberg へ戻せます。
ただし、HDFS と S3 の Iceberg warehouse をファイルコピーするのではなく、Spark + Iceberg で次のようにテーブルとして読み書きします。
glue_prod.${GLUE_DATABASE}.syslog_iceberg
↓
Spark
↓
hive_prod.logs.syslog_iceberg
今回作ったものです。
/opt/iceberg/bin/restore_iceberg_from_s3.sh/opt/iceberg/bin/restore_logs_from_s3.sh-
CONFIRM_RESTORE=1を必須にした手動リストア手順 - リストア前後の件数比較
これで、AWS 側の分析用 Iceberg コピーを、特定日データの復旧にも使えるようになります。
参考
- Apache Iceberg Spark writes: https://iceberg.apache.org/docs/latest/spark-writes/
- Apache Iceberg Spark configuration: https://iceberg.apache.org/docs/latest/spark-configuration/
- Apache Iceberg AWS integration: https://iceberg.apache.org/docs/latest/aws/
- Amazon Athena Iceberg tables: https://docs.aws.amazon.com/athena/latest/ug/querying-iceberg.html