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?

列指向データ作成,ローテーションスクリプト作成してみた

0
Last updated at Posted at 2026-01-04

はじめに

これまで作成した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:flag_mm: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を参照してください。

0
0
7

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?