はじめに
これまで完成しているsyslog収集環境に対し、auth.log収集を追加してみました。
収集仕様
- rsyslogdで収集したauth系ログを専用ファイルに書き込む。
- fluentd/kafkaにて専用topicsを使用して収集する。
- HDFS収集先は専用のものを使用する。
1.各収集対象ホスト設定
sudo tee /etc/rsyslog.d/90-forward-to-relay.conf << 'EOF'
*.* @@log1:514
EOF
authlogを送りたくない場合は以下の設定を実施
sudo tee /etc/rsyslog.d/90-forward-to-relay.conf << 'EOF'
# auth / authpriv(auth.log, authpriv.log)は転送除外
if ($syslogfacility-text == 'auth' or $syslogfacility-text == 'authpriv') then {
stop
}
# それ以外はすべて relay に転送
*.* @@log1:514
EOF
syslogについて、特定の時刻のみ別のホストに送りたい場合は以下の設定を実施
sudo tee /etc/rsyslog.d/90-forward-to-relay.conf << 'EOF'
# auth/authpriv は常に log1 へ転送
auth,authpriv.* @@log1:514
# 02:00~04:59 は syslog を hdptest へ転送
if (
$syslogfacility-text != 'auth' and
$syslogfacility-text != 'authpriv' and
re_match($timegenerated, "T0[2-4]:| 0[2-4]:")
) then {
action(type="omfwd" target="hdptest" port="514" protocol="tcp")
stop
}
# それ以外の時間帯は syslog を log1 へ転送
*.*;auth,authpriv.none @@log1:514
EOF
sudo rsyslogd -N1
→文法に問題がないことを確認する。
sudo systemctl restart rsyslog
→syslog と auth.log が同じ TCP 514(rsyslogd) 経路で log1 に届く
2. syslogd/fluentd設定(log1)
2.1 rsyslog設定
sudo tee /etc/rsyslog.d/20-relay-format.conf << 'EOF'
template(name="RelayFmt" type="string"
string="%timegenerated:::date-rfc3339% src_host=%fromhost% program=%programname% msg=%msg%\n"
)
# auth/authpriv → auth専用
if ($syslogfacility-text == "auth" or $syslogfacility-text == "authpriv") then {
action(type="omfile" file="/var/log/relay_auth.log" template="RelayFmt")
stop
}
# それ以外 → syslog側
action(type="omfile" file="/var/log/relay_syslog.log" template="RelayFmt")
EOF
sudo tee /etc/logrotate.d/relay_authlog <<'EOF'
/var/log/relay_auth.log {
daily
rotate 14
compress
missingok
notifempty
copytruncate
}
EOF
sudo rsyslogd -N1
sudo systemctl restart rsyslog
2.2 動作確認
tail -f /var/log/authlog
tail -f /var/log/relay_authlog.log
→auth系/syslog系ログが蓄積されること。
テスト:
logger -p auth.info "authlog forward test from $(hostname) $(date -Is)"
→”sudo tail -n 5 /var/log/authlog”を入力して、上記ログが参照できること。
→”sudo tail -n 5 /var/log/relay_authlog.log”を入力して、
対象ホストからの上記ログが出てること。
2.3. program/tag で Kafka topic を分ける
sudo mkdir -p /var/log/fluent/buffer
sudo chown -R _fluentd:_fluentd /var/log/fluent
sudo chmod 0755 /var/log/fluent
sudo tee /etc/fluent/fluentd.conf <<'EOF'
@include /etc/fluent/conf.d/*.conf
EOF
sudo mkdir /etc/fluent/conf.d/
sudo tee /etc/fluent/conf.d/10-source-relay_syslog.conf <<'EOF'
# ファイルを整形してsyslog.relayタグつけの実施
<source>
@type tail
path /var/log/relay_syslog.log
pos_file /var/log/fluent/relay_syslog.pos
tag syslog.relay
read_from_head true
<parse>
@type regexp
expression /^(?<ts>\S+)\s+src_host=(?<src_host>\S+)\s+program=(?<program>\S+)\s+msg=(?<message>.*)$/
</parse>
</source>
EOF
sudo tee /etc/fluent/conf.d/11-source-relay_auth.conf <<'EOF'
# ファイルを整形してauthlog.relayタグつけの実施
<source>
@type tail
path /var/log/relay_auth.log
pos_file /var/log/fluent/relay_auth.pos
tag authlog.relay
read_from_head true
<parse>
@type regexp
expression /^(?<ts>\S+)\s+src_host=(?<src_host>\S+)\s+program=(?<program>\S+)\s+msg=(?<message>.*)$/
</parse>
</source>
EOF
sudo tee /etc/fluent/conf.d/20-match-kafka-syslog.conf <<'EOF'
# syslog.relay → Kafka(syslog)
<match syslog.relay>
@type kafka2
brokers kafka1:9092,kafka2:9092,kafka3:9092
default_topic syslog
<format>
@type json
</format>
required_acks -1
compression_codec gzip
</match>
EOF
sudo tee /etc/fluent/conf.d/21-match-kafka-authlog.conf <<'EOF'
# authlog.relay → Kafka(authlog)
<match authlog.relay>
@type kafka2
brokers kafka1:9092,kafka2:9092,kafka3:9092
default_topic authlog
<format>
@type json
</format>
required_acks -1
compression_codec gzip
</match>
EOF
sudo fluentd --dry-run -c /etc/fluent/fluentd.conf
sudo systemctl restart fluentd
- syslogが必要ない場合
以下のファイルは必要ない。
sudo rm -f /etc/fluent/conf.d/10-source-relay_syslog.conf
sudo rm -f /etc/fluent/conf.d/20-match-kafka-syslog.conf
- auth.logが必要ない場合
以下のファイルは必要ない。
sudo rm -f /etc/fluent/conf.d/11-source-relay_auth.conf
sudo rm -f /etc/fluent/conf.d/21-match-kafka-authlog.conf
3. Kafka設定(kafkaシェルを設置してるホストで実施)
3.1. トピック作成(パーティション設計)
/opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka1:9092 --create \
--topic authlog --partitions 3 --replication-factor 3
kafkaが1台の場合は”--replication-factor 1”に変更する。
3.2. topic設定後の確認
Kafka に入ってるか確認
- topic一覧
/opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka1:9092 --list | egrep 'authlog'
- オフセット(流入が増えていればOK)
/opt/kafka/bin/kafka-get-offsets.sh --bootstrap-server kafka1:9092 --topic authlog
- 流れてるデータの目視
以下のコマンドを打ってしばらく待ち、ログが流れてくることを確認する。
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server kafka1:9092 \
--topic authlog
4. HDFS準備(hdfsコマンドが打てるホストで実施。)
sudo -u hadoop hdfs dfs -mkdir -p /data/kafka/authlog
sudo -u hadoop hdfs dfs -chmod -R 775 /data/kafka/authlog
sudo -u hadoop hdfs dfs -chown -R kafka:hadoop /data/kafka/authlog
sudo -u hadoop hdfs dfs -mkdir -p /logs/authlog
sudo -u hadoop hdfs dfs -chmod 1777 /logs/authlog
5. HDFS Sink Connector作成(ホストはどこでも実施してもよい。)
authlog用のHDFS Sink Connector作成する。
cat > hdfs-sink-authlog.json <<'EOF'
{
"name": "hdfs-sink-authlog",
"config": {
"connector.class": "io.confluent.connect.hdfs.HdfsSinkConnector",
"tasks.max": "1",
"topics": "authlog",
"hdfs.url": "hdfs://cluster1",
"hadoop.conf.dir": "/etc/hadoop/conf",
"topics.dir": "/data/kafka",
"logs.dir": "/logs",
"format.class": "io.confluent.connect.hdfs.json.JsonFormat",
"partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
"path.format": "yyyy/MM/dd/HH",
"partition.duration.ms": "3600000",
"timezone": "Asia/Tokyo",
"locale": "en",
"timestamp.extractor": "Wallclock",
"flush.size": "500",
"rotate.interval.ms": "600000",
"rotate.schedule.interval.ms": "600000",
"schema.compatibility": "NONE"
}
}
EOF
curl -X POST http://connect1:8083/connectors \
-H "Content-Type: application/json" \
--data @hdfs-sink-authlog.json
curl -s http://connect1:8083/connectors/hdfs-sink-authlog/status | jq
→stateがRUNNINGになること。
6. HDFS 側:配置を確認(hdfsコマンドが打てるホストで実施。)
hdfs dfs -ls -R -h /data/kafka/authlog | head
7. HDFS上に解析用データフォルダを作成する。
sudo -u hadoop hdfs dfs -mkdir -p /data/kafka/authlog_curated
sudo -u hadoop hdfs dfs -chown -R hive:hadoop /data/kafka/authlog_curated
sudo -u hadoop hdfs dfs -chmod -R 775 /data/kafka/authlog_curated
8. 解析用テーブル作成&HDFSローテーションシェル作成
8.1. 解析用テーブル作成SQL
sudo tee /usr/local/share/authlog_init.sql <<'EOF'
-- DB
CREATE DATABASE IF NOT EXISTS authlogs;
-- RAW
-- HDFSに保存されたJSONを1行1レコードとして読み込む
CREATE EXTERNAL TABLE IF NOT EXISTS authlogs.authlog_raw_day (
line STRING
)
STORED AS TEXTFILE
LOCATION '/data/kafka/authlog';
-- Curated
-- 分析用Parquetテーブル
CREATE EXTERNAL TABLE IF NOT EXISTS authlogs.authlog_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/authlog_curated';
-- RAW JSON解析用VIEW
CREATE OR REPLACE VIEW authlogs.v_authlog_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 authlogs.authlog_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
8.2. 解析用データ作成SQL
sudo tee /usr/local/share/authlog_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 authlogs.authlog_raw_day
SET LOCATION '/data/kafka/authlog/${DAY}';
INSERT OVERWRITE TABLE authlogs.authlog_curated PARTITION (dt='${DT}')
SELECT
host,
ts,
ts_text,
ts_raw,
severity,
program,
msg
FROM authlogs.v_authlog_parsed_day
WHERE dt='${DT}'
AND host IS NOT NULL
AND ts IS NOT NULL
AND msg IS NOT NULL;
EOF
8.3. SQL稼働&HDFSローテーションシェル
必要に応じ、JDBC_URLの変更を実施する。
sudo tee /usr/local/sbin/authlog_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/authlog}"
CURATED_LOCATION="${CURATED_LOCATION:-/data/kafka/authlog_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:-2}" # 昨日、一昨日のデータ取得
LOCKFILE="${LOCKFILE:-/var/lock/authlog_curate_rotate.lock}"
DRY_RUN="${DRY_RUN:-0}"
NO_LOAD="${NO_LOAD:-0}"
INIT_SQL="${INIT_SQL:-/usr/local/share/authlog_init.sql}"
LOAD_SQL="${LOAD_SQL:-/usr/local/share/authlog_load_dt.sql}"
# どれが無ければ init するか(必要に応じて増やせます)
NEEDED_DB="${NEEDED_DB:-authlogs}"
NEEDED_TABLES=(
"authlogs.authlog_raw_day"
"authlogs.authlog_curated"
)
NEEDED_VIEWS=(
"authlogs.v_authlog_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/authlog_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 authlogs.authlog_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 authlogs.authlog_curated DROP IF EXISTS PARTITION (dt='${d}');"
else
run_sql "ALTER TABLE authlogs.authlog_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
8.4. SQL&シェル権限付与
sudo chmod 755 /usr/local/sbin/authlog_curate_rotate.sh
sudo chmod 644 /usr/local/share/authlog_init.sql /usr/local/share/authlog_load_dt.sql
9. 解析テーブル作成&挿入
authlog取り込み当日ではデータが入らないため、
翌日に実施する。
sudo /usr/local/sbin/authlog_curate_rotate.sh
既存の authlogs.authlog_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 authlogs.authlog_curated;
"
sudo -u hadoop hdfs dfs -rm -r -skipTrash /data/kafka/authlog_curated/* || true
beeline -u "${JDBC_URL}" -n hive -f /usr/local/share/authlog_init.sql
sudo REBUILD_LAST_N_DAYS=1 /usr/local/sbin/authlog_curate_rotate.sh
Table ... has 7 columns, but query has 5 columns が出る場合は、
/usr/local/share/authlog_load_dt.sql が ts_text / ts_raw 追加前の古い内容です。
この章の /usr/local/share/authlog_load_dt.sql 作成手順を再実行してから、
authlog_curate_rotate.sh を実行し直します。
10. cron設定(毎日 04:10 JST)
sudo tee /etc/cron.d/authlog_curate_rotate >/dev/null <<'EOF'
SHELL=/bin/bash
PATH=/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin
TZ=Asia/Tokyo
10 4 * * * root /usr/local/sbin/authlog_curate_rotate.sh >> /var/log/authlog_curate_rotate.log 2>&1
EOF
11. 解析テーブル参照確認
日付部分はシェル実施日の前日に設定する。
zeppelinで以下のSQLを実施する。
%hive
set mapreduce.map.memory.mb=512;
set mapreduce.reduce.memory.mb=512;
set yarn.app.mapreduce.am.resource.mb=512;
-- データベース確認(authlogsがあること。)
SHOW DATABASES;
-- テーブル確認(各種テーブルがあること。)
USE authlogs;
SHOW TABLES;
-- パーティション確認(昨日、一昨日のパーティションがあること。)
SHOW PARTITIONS authlogs.authlog_curated;
-- 特定パーティションのレコード数(適当な数が出ること。)
SELECT count(*) FROM authlogs.authlog_curated WHERE dt='2026-01-05';
-- 特定パーティション最新10レコード(レコードが出てくること。)
SELECT * FROM authlogs.authlog_curated
WHERE dt='2026-01-05'
ORDER BY ts DESC
LIMIT 10;
12. 解析用データ、テーブルクリア方法
諸事情で解析用データ、テーブルを手動クリアする場合、以下コマンドを実施すること。
基本的にはraw(生ログ)はスクリプト以外から操作しない。
12.1 特定日のみ削除
%hive
ALTER TABLE authlogs.authlog_curated
DROP IF EXISTS PARTITION (dt='2026-01-03');
sudo -u hadoop hdfs dfs -rm -r -skipTrash /data/kafka/authlog_curated/dt=2026-01-03
12.2 全削除
- Zeppelin
drop view logs.v_authlog_parsed_day;
drop table authlogs.authlog_curated;
drop table authlogs.authlog_raw_day;
drop database authlogs;
- ope1
sudo -u hadoop hdfs dfs -rm -r /data/kafka/authlog_curated/*
おわりに
以降は以下の記事を参考にテーブル情報を書き換えれば解析できるようになります。