はじめに
これまで作成したHDFSに対して、列指向データ作成,ローテーションスクリプトを設置した。
以下仕様とする。
| 仕様項目 | パラメータ |
|---|---|
| 元データ保存期間 | 7日 |
| 列指向データ保持期間 | 14日 |
元データの削除は対象のデータをHDFSから直接削除
解析データについてはHDFS/Hiveメタデータから削除する。
0. 前提
以下の記事を基に運用ホスト(ope1)を構築していること。
シングルホストを構築している場合は、同ホストを使用して良い(以下の手順は不要)
これまで作成したsyslog関連テーブルなど削除すること。
コマンド:
- Zeppelin
drop view logs.v_syslog_parsed;
drop table logs.syslog;
drop database logs;
-- データベース確認
SHOW DATABASES;
-- テーブル確認
USE logs;
SHOW TABLES;
→データベース(logs)、テーブルが残っている場合(logs内に存在するテーブル)、適宜削除してください。
解析データを置くディレクトリを作成してください。(ope1)
sudo -u hadoop hdfs dfs -mkdir -p /data/kafka/syslog_curated
sudo -u hadoop hdfs dfs -chown -R hive:hadoop /data/kafka/syslog_curated
sudo -u hadoop hdfs dfs -chmod -R 775 /data/kafka/syslog_curated
1. 初期化SQL(1回だけ実行)
/usr/local/share/syslog_init.sql として保存
あらかじめZeppelin %hive で流してもよい。
sudo tee /usr/local/share/syslog_init.sql <<'EOF'
-- DB
CREATE DATABASE IF NOT EXISTS logs;
-- RAW
-- HDFSに保存されたJSONを1行1レコードとして読み込む
CREATE EXTERNAL TABLE IF NOT EXISTS logs.syslog_raw_day (
line STRING
)
STORED AS TEXTFILE
LOCATION '/data/kafka/syslog';
-- Curated
-- 分析用Parquetテーブル
CREATE EXTERNAL TABLE IF NOT EXISTS logs.syslog_curated (
host STRING,
ts TIMESTAMP,
ts_text STRING,
ts_raw STRING,
severity INT,
program STRING,
msg STRING
)
PARTITIONED BY (dt STRING)
STORED AS PARQUET
LOCATION '/data/kafka/syslog_curated';
-- RAW JSON解析用VIEW
CREATE OR REPLACE VIEW logs.v_syslog_parsed_day AS
WITH parsed AS (
SELECT
get_json_object(line, '$.src_host') AS host,
get_json_object(line, '$.ts') AS ts_raw,
get_json_object(line, '$.program') AS program,
get_json_object(line, '$.message') AS msg
FROM logs.syslog_raw_day
WHERE line IS NOT NULL
AND length(trim(line)) > 0
)
SELECT
host,
-- Hive 内の確認・集計用。Spark/Iceberg への受け渡しには ts_text を使う。
CAST(
regexp_replace(
substr(ts_raw, 1, 19),
'T',
' '
) AS TIMESTAMP
) AS ts,
regexp_replace(substr(ts_raw, 1, 19), 'T', ' ') AS ts_text,
ts_raw,
substr(ts_raw, 1, 10) AS dt,
program,
msg,
CASE
WHEN msg RLIKE '(?i)\\b(emerg|emergency)\\b' THEN 0
WHEN msg RLIKE '(?i)\\balert\\b' THEN 1
WHEN msg RLIKE '(?i)\\bcrit(ical)?\\b' THEN 2
WHEN msg RLIKE '(?i)\\b(err|error|failed|failure)\\b' THEN 3
WHEN msg RLIKE '(?i)\\bwarn(ing)?\\b' THEN 4
WHEN msg RLIKE '(?i)\\bnotice\\b' THEN 5
WHEN msg RLIKE '(?i)\\b(info|started|stopped)\\b' THEN 6
WHEN msg RLIKE '(?i)\\bdebug\\b' THEN 7
ELSE 6
END AS severity
FROM parsed
WHERE host IS NOT NULL
AND ts_raw IS NOT NULL
AND length(ts_raw) >= 19
AND msg IS NOT NULL;
EOF
2. 日次投入SQL(任意日 dt を流し込む)
これは dtを引数で差し替える前提です(シェルで置換)。
あらかじめ${DT}部分をyyyy-mm-ddで書き換え、
合わせて${DAY}部分もyyyy/mm/ddで書き換え、
Zeppelin %hive で流してもよい。
/usr/local/share/syslog_load_dt.sql として保存:
sudo tee /usr/local/share/syslog_load_dt.sql <<'EOF'
set mapreduce.map.memory.mb=512;
set mapreduce.reduce.memory.mb=512;
set yarn.app.mapreduce.am.resource.mb=512;
SET hive.exec.dynamic.partition=true;
SET hive.exec.dynamic.partition.mode=nonstrict;
set hive.mapred.supports.subdirectories=true;
set mapreduce.input.fileinputformat.input.dir.recursive=true;
ALTER TABLE logs.syslog_raw_day
SET LOCATION '/data/kafka/syslog/${DAY}';
INSERT OVERWRITE TABLE logs.syslog_curated PARTITION (dt='${DT}')
SELECT
host,
ts,
ts_text,
ts_raw,
severity,
program,
msg
FROM logs.v_syslog_parsed_day
WHERE dt='${DT}'
AND host IS NOT NULL
AND ts IS NOT NULL
AND msg IS NOT NULL;
EOF
3. シェル作成(初期化もできる・日次投入・ローテ実施)
3.1. 仕様
--init で初期化SQLを流す(初期化テーブルがない場合も流れる。)
通常実行は「昨日」投入(遅延到着対策)
curated(解析データ) 30日ローテ:DROP + HDFS実体削除
raw(生ログデータ) 7日ローテ:HDFS内の日ディレクトリ削除
多重起動防止:flock
DRY_RUN=1 既存データ削除停止
NO_LOAD=1 データインポート停止
3.2. シェル設置
/usr/local/sbin/syslog_curate_rotate.sh として保存
sudo tee /usr/local/sbin/syslog_curate_rotate.sh >/dev/null <<'EOF'
#!/usr/bin/env bash
set -euo pipefail
source /etc/profile.d/hadoop_hive_kafka.sh >/dev/null 2>&1 || true
export TZ="${TZ:-Asia/Tokyo}"
export LC_ALL=C
# =====================
# 非致命エラー管理(Rundeckアラート抑止用)
# =====================
NONFATAL_ERRORS=0
warn_nonfatal() {
echo "[WARN] $*" >&2
NONFATAL_ERRORS=$((NONFATAL_ERRORS+1))
}
# ====== 設定 ======
BEELINE_BIN="${BEELINE_BIN:-/usr/bin/beeline}"
JDBC_URL="${JDBC_URL:-jdbc:hive2://master1:2181,master2:2181,master3:2181/;serviceDiscoveryMode=zooKeeper;zooKeeperNamespace=hiveserver2}"
HIVE_USER="${HIVE_USER:-hive}"
RAW_BASE="${RAW_BASE:-/data/kafka/syslog}"
CURATED_LOCATION="${CURATED_LOCATION:-/data/kafka/syslog_curated}"
RAW_RETENTION_DAYS="${RAW_RETENTION_DAYS:-7}" # 7日前元データ削除
CURATED_RETENTION_DAYS="${CURATED_RETENTION_DAYS:-30}" # 30日前解析データ削除
REBUILD_LAST_N_DAYS="${REBUILD_LAST_N_DAYS:-1}" # 昨日のデータ取得
LOCKFILE="${LOCKFILE:-/var/lock/syslog_curate_rotate.lock}"
DRY_RUN="${DRY_RUN:-0}"
NO_LOAD="${NO_LOAD:-0}"
INIT_SQL="${INIT_SQL:-/usr/local/share/syslog_init.sql}"
LOAD_SQL="${LOAD_SQL:-/usr/local/share/syslog_load_dt.sql}"
# どれが無ければ init するか(必要に応じて増やせます)
NEEDED_DB="${NEEDED_DB:-logs}"
NEEDED_TABLES=(
"logs.syslog_raw_day"
"logs.syslog_curated"
)
NEEDED_VIEWS=(
"logs.v_syslog_parsed_day"
)
need_cmd() { command -v "$1" >/dev/null 2>&1 || { echo "[ERROR] command not found: $1" >&2; exit 1; }; }
run_sql() {
local sql="$1"
"${BEELINE_BIN}" -u "${JDBC_URL}" -n "${HIVE_USER}" \
--silent=true --showHeader=false --outputformat=tsv2 \
-e "${sql}"
}
run_sql_file_with_dt() {
local dt="$1"
local day="$2"
local tmp="/tmp/syslog_load_${dt}.sql"
sed -e "s/\${DT}/${dt}/g" -e "s|\${DAY}|${day}|g" "${LOAD_SQL}" > "${tmp}"
"${BEELINE_BIN}" -u "${JDBC_URL}" -n "${HIVE_USER}" \
--silent=true --showHeader=false --outputformat=tsv2 \
-f "${tmp}"
rm -f "${tmp}"
}
# ===== HDFS wrapper(ls/test/rm すべて hadoop ユーザーに統一)=====
hdfs_cmd() { sudo -n -u hadoop hdfs dfs "$@"; }
hdfs_rm_dir() {
local path="$1"
if [[ "${DRY_RUN}" == "1" ]]; then
echo "[DRY] hdfs dfs -rm -r -skipTrash ${path}"
return 0
fi
# ここが今回のキモ:
# rm失敗で set -e が発火してジョブ全体が exit 1 にならないようにする
local out rc
out="$(hdfs_cmd -rm -r -skipTrash "${path}" 2>&1)" || rc=$?
if [[ -z "${rc:-}" ]]; then
return 0
fi
# “既に無い”系は正常扱い
if grep -qiE 'No such file|does not exist|cannot find|File .* does not exist' <<<"${out}"; then
echo "[INFO] already gone: ${path}"
return 0
fi
warn_nonfatal "hdfs delete failed(rc=${rc}): ${path} :: ${out}"
return 0
}
hdfs_list_paths() {
# パスは最終列($NF)を使うのが安全
hdfs_cmd -ls "$1" 2>/dev/null | awk 'NF>=8{print $NF}'
}
# 0=empty, 1=not empty, 2=cannot judge(list不可)
hdfs_is_empty_dir() {
local dir="$1"
local out
if ! out="$(hdfs_list_paths "$dir")"; then
return 2
fi
[[ -z "$out" ]] && return 0 || return 1
}
# --- Hive存在チェック ---
hive_db_exists() {
local db="$1"
run_sql "SHOW DATABASES LIKE '${db}';" | awk '{print $1}' | grep -Fxq "${db}"
}
hive_table_exists() {
local fqn="$1" # db.table
local db="${fqn%%.*}"
local tbl="${fqn#*.}"
run_sql "SHOW TABLES IN ${db} LIKE '${tbl}';" | awk '{print $1}' | grep -Fxq "${tbl}"
}
hive_view_exists() {
local fqn="$1" # db.view
local db="${fqn%%.*}"
local vw="${fqn#*.}"
# 1) まず名前として存在するか確認
# beeline の表形式出力(| や空白)を吸収
if ! run_sql "SHOW TABLES IN ${db} LIKE '${vw}';" \
| sed 's/|/ /g' \
| sed 's/\r//g' \
| awk '{$1=$1; print}' \
| grep -Fxq "${vw}"; then
return 1
fi
# 2) DESCRIBE FORMATTED から Table Type を確認
# 出力形式差異:
# Table Type VIRTUAL_VIEW
# Table Type: VIRTUAL_VIEW
# | Table Type | VIRTUAL_VIEW |
# を吸収する
run_sql "DESCRIBE FORMATTED ${db}.${vw};" \
| sed 's/|/ /g' \
| sed 's/\r//g' \
| tr -s '[:space:]' ' ' \
| grep -Eq 'Table Type ?:? ?VIRTUAL_VIEW'
}
run_init_sql() {
echo "[INFO] running init sql: ${INIT_SQL}"
"${BEELINE_BIN}" -u "${JDBC_URL}" -n "${HIVE_USER}" \
--silent=true --showHeader=false --outputformat=tsv2 \
-f "${INIT_SQL}"
echo "[INFO] init done."
}
ensure_hive_objects() {
local need_init=0
echo "[INFO] checking Hive objects exist..."
if ! hive_db_exists "${NEEDED_DB}"; then
echo "[WARN] missing database: ${NEEDED_DB}"
need_init=1
fi
for t in "${NEEDED_TABLES[@]}"; do
if ! hive_table_exists "${t}"; then
echo "[WARN] missing table: ${t}"
need_init=1
fi
done
for v in "${NEEDED_VIEWS[@]}"; do
if ! hive_view_exists "${v}"; then
echo "[WARN] missing view: ${v}"
need_init=1
fi
done
if [[ "${need_init}" == "1" ]]; then
run_init_sql
else
echo "[INFO] Hive objects OK (no init needed)."
fi
}
usage() {
cat <<EOU
Usage:
$0 [--init]
Behavior:
- If required Hive objects are missing, init runs automatically.
- With --init, init runs regardless (VIEW will be replaced).
Env:
JDBC_URL, HIVE_USER, RAW_BASE, CURATED_LOCATION
RAW_RETENTION_DAYS, CURATED_RETENTION_DAYS, REBUILD_LAST_N_DAYS
DRY_RUN=1 (do not delete)
NO_LOAD=1 (skip load)
EOU
}
# ====== 引数 ======
FORCE_INIT=0
if [[ "${1:-}" == "--init" ]]; then
FORCE_INIT=1
elif [[ "${1:-}" != "" ]]; then
usage; exit 1
fi
# ====== 事前チェック ======
need_cmd flock
need_cmd date
need_cmd sed
need_cmd hdfs
need_cmd awk
need_cmd grep
need_cmd sudo
exec 9>"${LOCKFILE}"
flock -n 9 || { echo "[WARN] already running: ${LOCKFILE}"; exit 0; }
echo "[INFO] start: $(date '+%F %T %Z')"
echo "[INFO] JDBC_URL=${JDBC_URL} HIVE_USER=${HIVE_USER}"
echo "[INFO] RAW_BASE=${RAW_BASE} CURATED_LOCATION=${CURATED_LOCATION}"
echo "[INFO] retention: raw=${RAW_RETENTION_DAYS}d curated=${CURATED_RETENTION_DAYS}d rebuild_last=${REBUILD_LAST_N_DAYS}d DRY_RUN=${DRY_RUN}"
echo "[INFO] NO_LOAD=${NO_LOAD}"
echo "[INFO] checking Hive connection..."
run_sql "SELECT 1;" >/dev/null
echo "[INFO] checking HDFS access (as hadoop user)..."
hdfs_cmd -ls / >/dev/null
# ====== 初期化(冪等):無ければ自動init/--initなら強制init ======
if [[ "${FORCE_INIT}" == "1" ]]; then
run_init_sql
echo "[INFO] init-only mode (--init). exiting."
exit 0
else
ensure_hive_objects
fi
# ====== 直近N日を投入(OVERWRITE)=====
# DRY_RUN=1 のときに投入も止めたい場合は NO_LOAD=1 を使う
if [[ "${NO_LOAD}" == "1" ]]; then
echo "[INFO] skip load (NO_LOAD=1)"
else
echo "[INFO] load last ${REBUILD_LAST_N_DAYS} day(s) ..."
for k in $(seq 1 "${REBUILD_LAST_N_DAYS}"); do
DT="$(date -d "${k} day ago" +%F)"
DAY="$(date -d "${k} day ago" +%Y/%m/%d)"
echo "[INFO] - load dt=${DT} day=${DAY}"
run_sql_file_with_dt "${DT}" "${DAY}"
done
fi
# ==============================
# curated ローテ(非致命にしたい区間)
# ==============================
rotate_curated() {
# ====== curated ローテ(メタストア掃除 + HDFS実体掃除:実在パーティションのみ)=====
CUTOFF="$(date -d "${CURATED_RETENTION_DAYS} days ago" +%F)"
echo "[INFO] rotate curated partitions where dt < ${CUTOFF} (from metastore) ..."
PARTS="$(run_sql "SHOW PARTITIONS logs.syslog_curated;" 2>/dev/null || true)"
if [[ -z "${PARTS}" ]]; then
echo "[INFO] no partitions found (or cannot list partitions)."
else
echo "${PARTS}" \
| awk -F= '/^dt=/{print $2}' \
| sort \
| while read -r d; do
[[ -z "${d}" ]] && continue
# 文字列比較(YYYY-MM-DDなので辞書順=日付順)
if [[ "${d}" < "${CUTOFF}" ]]; then
echo "[INFO] - drop curated dt=${d}"
# 1) メタストア(パーティション)削除
if [[ "${DRY_RUN}" == "1" ]]; then
echo "[DRY] ALTER TABLE logs.syslog_curated DROP IF EXISTS PARTITION (dt='${d}');"
else
run_sql "ALTER TABLE logs.syslog_curated DROP IF EXISTS PARTITION (dt='${d}');"
fi
# 2) HDFS実体削除(外部テーブルなので明示的に消す)
curated_path="${CURATED_LOCATION}/dt=${d}"
if hdfs_cmd -test -d "${curated_path}" 2>/dev/null; then
echo "[INFO] - delete curated data ${curated_path}"
hdfs_rm_dir "${curated_path}"
fi
fi
done
fi
return 0
}
# ==============================
# raw ローテ(非致命にしたい区間)
# ==============================
rotate_raw() {
# ====== raw ローテ(HDFS実体のみ)=====
RAW_CUTOFF="$(date -d "${RAW_RETENTION_DAYS} days ago" +%F)"
CUTOFF_YM="${RAW_CUTOFF:0:7}"
TODAY="$(date +%F)"
echo "[INFO] rotate raw where date < ${RAW_CUTOFF} (BASE=${RAW_BASE}) ..."
# 事故防止:future cutoff を拒否(全消し防止) ←ここは致命でOK
if [[ "${RAW_CUTOFF}" > "${TODAY}" ]]; then
echo "[ERROR] RAW_CUTOFF is in the future. today=${TODAY} RAW_CUTOFF=${RAW_CUTOFF}" >&2
return 2
fi
# RAW_BASEがlistできないならここで止める(“見えないのに消す”事故防止) ←ここも致命でOK
if ! hdfs_cmd -ls "${RAW_BASE}" >/dev/null 2>&1; then
echo "[ERROR] cannot list RAW_BASE as hadoop user: ${RAW_BASE}" >&2
return 1
fi
# ここは “たまにlsが失敗する” だけで set -e が発火しがちなので保険を入れる
local ylist
if ! ylist="$(hdfs_list_paths "${RAW_BASE}")"; then
warn_nonfatal "cannot enumerate RAW_BASE children: ${RAW_BASE} (skip raw rotate)"
return 0
fi
echo "${ylist}" \
| awk '$0 ~ /\/[0-9]{4}$/ {print}' \
| while read -r ydir; do
mlist="$(hdfs_list_paths "${ydir}" || true)"
if [[ -z "${mlist}" ]]; then
# listできない場合は削除しない(安全側)
if hdfs_is_empty_dir "${ydir}"; then
echo "[INFO] - delete empty year ${ydir}"
hdfs_rm_dir "${ydir}"
else
echo "[WARN] cannot judge year dir ${ydir}; skip deleting (safety)"
fi
continue
fi
echo "${mlist}" \
| awk '$0 ~ /\/[0-9]{2}$/ {print}' \
| while read -r mdir; do
y="$(echo "${mdir}" | awk -F/ '{print $(NF-1)}')" # YYYY
m="$(echo "${mdir}" | awk -F/ '{print $NF}')" # MM
ym="${y}-${m}"
if [[ ! "${y}" =~ ^[0-9]{4}$ || ! "${m}" =~ ^[0-9]{2}$ ]]; then
echo "[WARN] skip non ym dir: ${mdir}"
continue
fi
# --- 古い月は月ごと削除(元の仕様のまま) ---
if [[ "${ym}" < "${CUTOFF_YM}" ]]; then
echo "[INFO] - delete raw month ${mdir} (ym=${ym})"
hdfs_rm_dir "${mdir}"
continue
fi
# cutoff月より新しい月は触らない
if [[ "${ym}" > "${CUTOFF_YM}" ]]; then
continue
fi
# --- cutoff月だけ:日単位で削除 ---
dlist="$(hdfs_list_paths "${mdir}" || true)"
if [[ -z "${dlist}" ]]; then
if hdfs_is_empty_dir "${mdir}"; then
echo "[INFO] - delete empty month ${mdir}"
hdfs_rm_dir "${mdir}"
else
echo "[WARN] cannot judge month dir ${mdir}; skip deleting (safety)"
fi
continue
fi
echo "${dlist}" \
| awk '$0 ~ /\/[0-9]{2}$/ {print}' \
| while read -r ddir; do
d="$(echo "${ddir}" | awk -F/ '{print $NF}')" # DD
if [[ ! "${d}" =~ ^[0-9]{2}$ ]]; then
echo "[WARN] skip non date dir: ${ddir}"
continue
fi
ymd="${y}-${m}-${d}"
if [[ "${ymd}" < "${RAW_CUTOFF}" ]]; then
echo "[INFO] - delete raw day ${ddir} (date=${ymd})"
hdfs_rm_dir "${ddir}"
fi
done
done
# 年ディレクトリが空になったら削除(list不能なら削除しない)
if hdfs_is_empty_dir "${ydir}"; then
echo "[INFO] - delete empty year ${ydir}"
hdfs_rm_dir "${ydir}"
else
rc=$?
[[ "$rc" -eq 2 ]] && echo "[WARN] cannot judge year dir ${ydir}; skip deleting (safety)"
fi
done
return 0
}
# ==============================
# ローテ実行:ここだけ「落ちない」ようにする
# ==============================
# ※ただし rotate_raw 内の致命条件(future cutoff / RAW_BASE list不可)は return 1 で伝える
set +e
rotate_curated
rc1=$?
rotate_raw
rc2=$?
set -e
# raw側で致命(事故防止)を踏んだ場合だけは、従来どおりジョブ失敗にする
if [[ "${rc1}" -ne 0 ]]; then
warn_nonfatal "rotate_curated returned rc=${rc1}"
fi
if [[ "${rc2}" -eq 2 ]]; then
echo "[ERROR] rotate_raw fatal guard triggered (future cutoff)" >&2
exit 1
elif [[ "${rc2}" -ne 0 ]]; then
warn_nonfatal "rotate_raw returned rc=${rc2} (ignored)"
fi
if [[ "${NONFATAL_ERRORS}" -gt 0 ]]; then
echo "[WARN] finished with non-fatal errors: ${NONFATAL_ERRORS}"
fi
echo "[INFO] done: $(date '+%F %T %Z')"
exit 0
EOF
3.3. シェル権限設定
sudo chmod 755 /usr/local/sbin/syslog_curate_rotate.sh
sudo chmod 644 /usr/local/share/syslog_init.sql /usr/local/share/syslog_load_dt.sql
4. 実行手順
4.1. 初期化(テーブル作成のみ)
sudo /usr/local/sbin/syslog_curate_rotate.sh --init
4.2. ドライラン(データインサートなし、データ削除なし)
初期化実施していない場合、自動的に初期化から稼働されます。
前日のデータは挿入されません。
削除対象がある場合返答内容で出てきます。
sudo DRY_RUN=1 NO_LOAD=1 /usr/local/sbin/syslog_curate_rotate.sh
4.3. ドライラン(データインサートのみ、データ削除なし)
初期化実施していない場合、自動的に初期化から稼働されます。
前日のデータが挿入されます。
削除対象がある場合返答内容で出てきます。
sudo DRY_RUN=1 /usr/local/sbin/syslog_curate_rotate.sh
4.4. 本番実行
初期化実施していない場合、自動的に初期化から稼働されます。
一昨日,前日のデータが挿入されます。
削除対象がある場合削除されます。(ローテーション実施)
sudo /usr/local/sbin/syslog_curate_rotate.sh
既存の logs.syslog_curated を作成済みの場合、CREATE TABLE IF NOT EXISTS では
ts_text / ts_raw は追加されません。
時刻解釈を揃える目的では、テーブル定義を作り直して 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
5. cron設定(毎日 03:00 JST)
sudo tee /etc/cron.d/syslog_curate_rotate >/dev/null <<'EOF'
SHELL=/bin/bash
PATH=/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin
TZ=Asia/Tokyo
00 3 * * * root /usr/local/sbin/syslog_curate_rotate.sh >> /var/log/syslog_curate_rotate.log 2>&1
EOF
6. 動作確認コマンド
zeppelinで実施
上記シェル実施後に確認してもよい。
日付は適当な日時にすること。
%hive
set mapreduce.map.memory.mb=512;
set mapreduce.reduce.memory.mb=512;
set yarn.app.mapreduce.am.resource.mb=512;
-- データベース確認(logsがあること。)
SHOW DATABASES;
-- テーブル確認(各種テーブルがあること。)
USE logs;
SHOW TABLES;
-- パーティション確認(昨日、一昨日のパーティションがあること。)
SHOW PARTITIONS logs.syslog_curated;
-- 特定パーティションのレコード数(適当な数が出ること。)
SELECT count(*) FROM logs.syslog_curated WHERE dt='2026-01-04';
-- 特定パーティション最新10レコード(レコードが出てくること。)
SELECT * FROM logs.syslog_curated
WHERE dt='2026-01-04'
ORDER BY ts DESC
LIMIT 10;
7. 解析方法
今後は以下のサンプルを参考に解析を行ってください。
7.1. 日単位グラフ
X軸:yyyy-MM-dd
%hive
set mapreduce.map.memory.mb=512;
set mapreduce.reduce.memory.mb=512;
set yarn.app.mapreduce.am.resource.mb=512;
SELECT
host,
dt AS x_day,
COUNT(*) AS cnt
FROM logs.syslog_curated
WHERE dt BETWEEN '2026-01-03' AND '2026-01-04'
AND severity <= 4
GROUP BY
host,
dt
ORDER BY
x_day,
host;
7.2. 時単位グラフ
X軸:yyyy-MM-dd HH
%hive
set mapreduce.map.memory.mb=512;
set mapreduce.reduce.memory.mb=512;
set yarn.app.mapreduce.am.resource.mb=512;
SELECT
host,
concat(dt, ' ', lpad(hour(ts), 2, '0')) AS x_time,
COUNT(*) AS cnt
FROM logs.syslog_curated
WHERE dt BETWEEN '2026-01-03' AND '2026-01-04'
AND severity <= 4
GROUP BY
host,
concat(dt, ' ', lpad(hour(ts), 2, '0'))
ORDER BY
x_time,
host;
7.3. 分単位グラフ
X軸:yyyy-MM-dd HH:mm
%hive
set mapreduce.map.memory.mb=512;
set mapreduce.reduce.memory.mb=512;
set yarn.app.mapreduce.am.resource.mb=512;
SELECT
host,
concat(
dt, ' ',
lpad(hour(ts), 2, '0'), ':',
lpad(minute(ts), 2, '0')
) AS x_time,
COUNT(*) AS cnt
FROM logs.syslog_curated
WHERE dt = '2026-01-04'
AND ts >= cast('2026-01-04 21:00:00' as timestamp)
AND ts < cast('2026-01-04 23:00:00' as timestamp)
AND severity <= 4
GROUP BY
host,
concat(
dt, ' ',
lpad(hour(ts), 2, '0'), ':',
lpad(minute(ts), 2, '0')
)
ORDER BY
x_time,
host;
7.4. 秒単位グラフ(second)
X軸:yyyy-MM-dd HH
ss
%hive
set mapreduce.map.memory.mb=512;
set mapreduce.reduce.memory.mb=512;
set yarn.app.mapreduce.am.resource.mb=512;
SELECT
host,
cast(concat(
dt,' ',
lpad(hour(ts),2,'0'), ':',
lpad(minute(ts),2,'0'), ':',
lpad(second(ts),2,'0')
) as timestamp) AS x_time,
count(*) AS cnt
FROM logs.syslog_curated
WHERE dt='2026-01-04'
AND ts >= cast('2026-01-04 23:00:00' as timestamp)
AND ts < cast('2026-01-04 23:10:00' as timestamp)
GROUP BY host, cast(concat(
dt,' ',
lpad(hour(ts),2,'0'), ':',
lpad(minute(ts),2,'0'), ':',
lpad(second(ts),2,'0')
) as timestamp)
ORDER BY x_time, host;
8. 解析用データ、テーブルクリア方法
諸事情で解析用データ、テーブルを手動クリアする場合、以下コマンドを実施すること。
基本的にはraw(生ログ)はスクリプト以外から操作しない。
8.1 特定日のみ削除
%hive
ALTER TABLE logs.syslog_curated
DROP IF EXISTS PARTITION (dt='2026-01-03');
sudo -u hadoop hdfs dfs -rm -r -skipTrash /data/kafka/syslog_curated/dt=2026-01-03
8.2 全削除
- Zeppelin
drop view logs.v_syslog_parsed_day;
drop table logs.syslog_curated;
drop table logs.syslog_raw_day;
drop database logs;
- ope1
sudo -u hadoop hdfs dfs -rm -r /data/kafka/syslog_curated/*
9. HDFS取り込み前確認シェル
本シェルについて、HDFS内に壊れた生ログがあると以降全てコケてしまうため、
これも導入して取り込みシェル実施前に動かしてください。
9.1. シェル作成
9.1.1. 仕様
前日のフォルダ毎のものを確認する。
問題あるファイルがあったら退避する。
ログファイルに内容の情報を残す。
環境依存の箇所は変数化し、状況によって個別指定できるようにする。
9.1.2. シェル設置
/usr/local/sbin/hdfs_check.sh として保存
sudo tee /usr/local/sbin/hdfs_check.sh >/dev/null <<'EOF'
#!/bin/bash
set -euo pipefail
# 例: hdfs://hdptest:9000 / hdfs://hdptest2:9000 / hdfs://cluster1
HDFS_URI="${HDFS_URI:-hdfs://cluster1}"
# 既定は syslog(authlog を見たい場合は env で上書き)
BASE_ROOT="${BASE_ROOT:-/data/kafka/syslog}"
QUARANTINE_ROOT="${QUARANTINE_ROOT:-/data/kafka/_corrupt/syslog}"
LOGFILE="${LOGFILE:-/var/log/hdfs_check.log}"
# 昨日分をチェック(例: 2026/02/21)
YESTERDAY="$(date -d 'yesterday' '+%Y/%m/%d')"
CHECK_DIR="${BASE_ROOT}/${YESTERDAY}"
log(){
local now
now="$(date '+%Y-%m-%d %H:%M:%S')"
echo "[$now] $*" | tee -a "$LOGFILE"
}
# fsck は -fs が使えないのでフルURIで
FSCK_TARGET="${HDFS_URI}${CHECK_DIR}"
log "Start: fs=${HDFS_URI} target=${CHECK_DIR}"
# 対象ディレクトリ存在確認(dfs は -fs が使える)
if ! sudo -n -u hadoop hdfs dfs -fs "$HDFS_URI" -test -d "$CHECK_DIR"; then
log "Directory not found: $CHECK_DIR"
exit 0
fi
# CORRUPT 抽出
# - fsck 出力の「/path: ... CORRUPT ...」形式の行だけ拾う(サマリ行の "CORRUPT" を除外)
# - 1列目 "/path:" の末尾 ":" を外してパス化
mapfile -t CORRUPT_FILES < <(
sudo -n -u hadoop hdfs fsck "$FSCK_TARGET" -files -blocks \
| awk '
$1 ~ /^\/.*:$/ && /CORRUPT/ {
sub(/:$/, "", $1);
print $1
}
' \
| sort -u
)
if [[ "${#CORRUPT_FILES[@]}" -eq 0 ]]; then
log "No corrupt files found"
exit 0
fi
log "Found ${#CORRUPT_FILES[@]} corrupt files"
for FILE in "${CORRUPT_FILES[@]}"; do
# 念のため:パスっぽくないトークンは捨てる(事故防止)
if [[ "$FILE" != /* && "$FILE" != hdfs://* ]]; then
log "Skip non-path token: $FILE"
continue
fi
# FILE は "/path" が基本だが、hdfs://... で来た場合に備えてパス部分へ寄せる
FILE_PATH="$FILE"
if [[ "$FILE_PATH" == hdfs://* ]]; then
FILE_PATH="${FILE_PATH#${HDFS_URI}}"
fi
# さらに保険:mv できるのは絶対パスだけ
if [[ "$FILE_PATH" != /* ]]; then
log "Skip invalid path after normalize: raw=$FILE normalized=$FILE_PATH"
continue
fi
REL_PATH="${FILE_PATH#${BASE_ROOT}/}"
DEST_DIR="$(dirname "${QUARANTINE_ROOT}/${REL_PATH}")"
log "Quarantine: $FILE_PATH -> $DEST_DIR"
sudo -n -u hadoop hdfs dfs -fs "$HDFS_URI" -mkdir -p "$DEST_DIR"
sudo -n -u hadoop hdfs dfs -fs "$HDFS_URI" -mv "$FILE_PATH" "$DEST_DIR/"
done
log "Completed"
EOF
9.1.3. シェル権限設定
sudo chmod 755 /usr/local/sbin/hdfs_check.sh
ログは以下の内容を参照する。
tail -f /var/log/hdfs_check.log
9.2. 実行手順
sudo /usr/local/sbin/hdfs_check.sh
- hdfsを個別に指定して、authlogを指定する場合
sudo HDFS_URI="hdfs://cluster1" BASE_ROOT=/data/kafka/authlog QUARANTINE_ROOT=/data/kafka/_corrupt/authlog /usr/local/sbin/hdfs_check.sh
9.3. cron設定(データローテーション実施前に行う。)
sudo tee /etc/cron.d/syslog_curate_rotate >/dev/null <<'EOF'
SHELL=/bin/bash
PATH=/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin
TZ=Asia/Tokyo
00 3 * * * root /usr/local/sbin/hdfs_check.sh 2>&1 && /usr/local/sbin/syslog_curate_rotate.sh >> /var/log/syslog_curate_rotate.log 2>&1
EOF
10. 運用上の補足
実運用では以下の構成を推奨します。
HDFS(JSON)
↓
Hive RAW
↓
Hive Curated
↓
Iceberg
↓
Trino
↓
Grafana
JSON解析はRAW→Curatedで実施し、その後はIcebergを参照してください。