0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

Hadoop収集対象にauth.log追加してみた

0
Last updated at Posted at 2026-01-05

はじめに

これまで完成している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.sqlts_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/*

おわりに

以降は以下の記事を参考にテーブル情報を書き換えれば解析できるようになります。

0
0
2

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?