はじめに
今まで構築してきたHadoop関連の集大成として、
各ホストのsyslogを収集してみました。
事前に以下記事でHadoop&Hive&Zeppelinを構築していること。
構成
パラメータ
OS: Ubuntu / Debian
Hadoop / HDFS / Hive / Zeppelin は 既に構築済み
HDFS URI: hdfs://cluster1
Kafka: 1台 PoC(KRaft)
syslog: /var/log/syslog
追加ホスト名
syslog-relay: log1
Kafka + Connect: kafka1
シングルノードに移し替えも可能で、ホスト名の部分をlocalhostに置き換えれば稼働可能。
1. syslog-relay 構築(rsyslog + Fluentd)
1-1. rsyslog 受信 + src_host を厳密に埋め込む
sudo tee /etc/rsyslog.d/10-listen.conf << 'EOF'
module(load="imtcp")
input(type="imtcp" port="514")
module(load="imudp")
input(type="imudp" port="514")
EOF
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"
)
action(type="omfile" file="/var/log/relay_syslog.log" template="RelayFmt")
EOF
sudo systemctl restart rsyslog
sudo ss -lntp | grep 514
sudo tee /etc/logrotate.d/relay_syslog <<'EOF'
/var/log/relay_syslog.log {
daily
rotate 14
compress
missingok
notifempty
copytruncate
}
EOF
1-2. Fluentd(td-agent)インストール
sudo apt update
sudo apt install -y curl gnupg
curl -fsSL https://packages.treasuredata.com/GPG-KEY-td-agent | sudo apt-key add -
echo "deb https://packages.treasuredata.com/5/ubuntu/$(lsb_release -sc)/ $(lsb_release -sc) contrib" \
| sudo tee /etc/apt/sources.list.d/td-agent.list
sudo apt update
sudo apt install -y td-agent
1-3. Fluentd 設定(RFC5424風パース)
sudo tee /etc/fluent/fluentd.conf << 'EOF'
<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>
<match syslog.relay>
@type kafka2
brokers kafka1:9092
default_topic syslog
<format>
@type json
</format>
required_acks 1
compression_codec gzip
</match>
EOF
sudo install -d -o _fluentd -g _fluentd -m 0755 /var/log/fluent
sudo chown _fluentd:_fluentd /var/log/fluent
sudo systemctl restart fluentd
2. 各ホスト → syslog-relay(rsyslog 転送)
syslog を出す全ホストで実行:
sudo tee /etc/rsyslog.d/90-forward-to-relay.conf << 'EOF'
*.* @@log1:514
EOF
sudo systemctl restart rsyslog
テスト:
logger -p user.info "syslog forward test from $(hostname) $(date -Is)"
→”sudo tail -n 5 /var/log/syslog”を入力して、上記ログが参照できること。
→”sudo tail -n 5 /var/log/relay_syslog.log”を入力して、
対象ホストからの上記ログが出てること。
3 Kafka + Kafka Connect 同居ホスト構築
3-1. Java & ユーザー
sudo apt update
sudo apt install -y openjdk-17-jre-headless curl tar unzip
sudo useradd --system --home /var/lib/kafka --shell /usr/sbin/nologin kafka || true
sudo mkdir -p /var/lib/kafka /opt/kafka-connect
sudo chown -R kafka:kafka /var/lib/kafka /opt/kafka-connect
3-2. Hadoop conf 配置
sudo mkdir -p /etc/hadoop/conf
sudo scp master1:/etc/hadoop/conf/{core-site.xml,hdfs-site.xml} /etc/hadoop/conf/
→下記記事の対象ファイルをコピペしてもよい。
https://qiita.com/naritomo08/items/7891e0040cac74b6475f
sudo chown -R kafka:kafka /etc/hadoop
3-3. Kafka 配置
cd /tmp
wget https://archive.apache.org/dist/kafka/3.7.1/kafka_2.13-3.7.1.tgz
tar xzf kafka_2.13-3.7.1.tgz
sudo mv kafka_2.13-3.7.1 /opt/kafka
sudo chown -R kafka:kafka /opt/kafka
cd
3-4. Kafka(KRaft)設定
sudo tee /opt/kafka/server.properties << 'EOF'
process.roles=broker,controller
node.id=1
controller.listener.names=CONTROLLER
controller.quorum.voters=1@127.0.0.1:9093
listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://127.0.0.1:9093
advertised.listeners=PLAINTEXT://kafka1:9092
listener.security.protocol.map=PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT
inter.broker.listener.name=PLAINTEXT
log.dirs=/var/lib/kafka/data
num.partitions=3
default.replication.factor=1
min.insync.replicas=1
offsets.topic.replication.factor=1
transaction.state.log.replication.factor=1
transaction.state.log.min.isr=1
auto.create.topics.enable=false
EOF
3-5. KRaft 初期化(※1回だけ)
sudo -u kafka mkdir -p /var/lib/kafka/data
CLUSTER_ID=$(/opt/kafka/bin/kafka-storage.sh random-uuid)
sudo -u kafka /opt/kafka/bin/kafka-storage.sh format \
-t "$CLUSTER_ID" \
-c /opt/kafka/server.properties
3-6. Kafka systemd
sudo tee /etc/systemd/system/kafka.service << 'EOF'
[Unit]
Description=Apache Kafka (KRaft)
After=network-online.target
[Service]
Type=simple
User=kafka
Environment="JAVA_HOME=/usr/lib/jvm/java-17-openjdk-amd64"
WorkingDirectory=/opt/kafka
ExecStart=/opt/kafka/bin/kafka-server-start.sh /opt/kafka/server.properties
Restart=on-failure
LimitNOFILE=100000
[Install]
WantedBy=multi-user.target
EOF
sudo systemctl daemon-reload
sudo systemctl enable --now kafka
3-7. Kafka Connect 設定
sudo tee /opt/kafka/config/connect-distributed.properties << 'EOF'
bootstrap.servers=kafka1:9092
group.id=connect-hdfs
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter.schemas.enable=false
config.storage.topic=connect-configs
offset.storage.topic=connect-offsets
status.storage.topic=connect-status
config.storage.replication.factor=1
offset.storage.replication.factor=1
status.storage.replication.factor=1
plugin.path=/opt/kafka-connect/plugins
rest.port=8083
EOF
3-8. Kafka Connect systemd
sudo tee /etc/systemd/system/kafka-connect.service << 'EOF'
[Unit]
Description=Kafka Connect
After=kafka.service
[Service]
Type=simple
User=kafka
Environment="JAVA_HOME=/usr/lib/jvm/java-17-openjdk-amd64"
Environment="HADOOP_CONF_DIR=/etc/hadoop/conf"
ExecStart=/opt/kafka/bin/connect-distributed.sh /opt/kafka/config/connect-distributed.properties
Restart=on-failure
[Install]
WantedBy=multi-user.target
EOF
sudo systemctl daemon-reload
sudo systemctl enable --now kafka-connect
3-8. Connectorログローテーション
ローテーションシェル作成(7日でローテーション実施)
sudo tee /usr/local/sbin/cleanup-kafka-connect-logs.sh >/dev/null <<'EOF'
#!/usr/bin/env bash
set -euo pipefail
LOG_DIR="/opt/kafka/logs"
# 24時間後に圧縮
COMPRESS_MINUTES="${COMPRESS_MINUTES:-1440}"
# 7日後に削除
RETENTION_MINUTES="${RETENTION_MINUTES:-10080}"
echo "[INFO] cleanup start"
echo "[INFO] log_dir=${LOG_DIR}"
echo "[INFO] compress_minutes=${COMPRESS_MINUTES}"
echo "[INFO] retention_minutes=${RETENTION_MINUTES}"
echo "[INFO] compressing old Kafka Connect logs"
find "${LOG_DIR}" \
-maxdepth 1 \
-type f \
-name 'connect.log.20*' \
! -name '*.gz' \
-mmin "+${COMPRESS_MINUTES}" \
-print0 |
xargs -0 -r -n 1 gzip --verbose
echo "[INFO] deleting expired compressed logs"
find "${LOG_DIR}" \
-maxdepth 1 \
-type f \
-name 'connect.log.20*.gz' \
-mmin "+${RETENTION_MINUTES}" \
-print \
-delete
echo "[INFO] cleanup completed"
EOF
sudo chmod 0755 /usr/local/sbin/cleanup-kafka-connect-logs.sh
ローテーションサービス作成
sudo tee /etc/systemd/system/kafka-connect-log-cleanup.service >/dev/null <<'EOF'
[Unit]
Description=Cleanup old Kafka Connect log files
[Service]
Type=oneshot
ExecStart=/usr/local/sbin/cleanup-kafka-connect-logs.sh
EOF
ローテーションtimer作成(毎日3:15稼働)
sudo tee /etc/systemd/system/kafka-connect-log-cleanup.timer >/dev/null <<'EOF'
[Unit]
Description=Daily Kafka Connect log cleanup
[Timer]
OnCalendar=*-*-* 03:15:00
Persistent=true
Unit=kafka-connect-log-cleanup.service
[Install]
WantedBy=timers.target
EOF
ローテーション有効化
sudo systemctl daemon-reload
sudo systemctl enable --now kafka-connect-log-cleanup.timer
systemctl list-timers kafka-connect-log-cleanup.timer
手動テスト
sudo systemctl start kafka-connect-log-cleanup.service
sudo journalctl -u kafka-connect-log-cleanup.service -n 50 --no-pager
4. Kafka → HDFS設定
kafka1で実施
/opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka1:9092 \
--create --topic syslog --partitions 3 --replication-factor 1
master1で実施
sudo -u hadoop hdfs dfs -mkdir -p /data/kafka/syslog
sudo -u hadoop hdfs dfs -chown -R kafka:hadoop /data/kafka
sudo -u hadoop hdfs dfs -chmod -R 775 /data/kafka
sudo -u hadoop hdfs dfs -mkdir -p /logs/syslog
sudo -u hadoop hdfs dfs -chmod 1777 /logs
sudo -u hadoop hdfs dfs -chmod 1777 /logs/syslog
5. Confluent プラグインを入れる
5-1. プラグイン置き場を作る
sudo mkdir -p /opt/kafka/plugins
sudo chown -R kafka:kafka /opt/kafka/plugins
5-2. Confluent HDFS Sink を入れる(Confluent Hub Client)
5-2-1. confluent-hub を入れる
cd /tmp
curl -fsSLO https://client.hub.confluent.io/confluent-hub-client-latest.tar.gz
mkdir -p /tmp/confluent-hub-client
tar -xzvf confluent-hub-client-latest.tar.gz -C /tmp/confluent-hub-client
sudo mkdir -p /opt/confluent-hub
sudo cp -a confluent-hub-client/* /opt/confluent-hub/
sudo ln -sf /opt/confluent-hub/bin/confluent-hub /usr/local/bin/confluent-hub
rm -rf confluent-hub-client*
cd
5-2-2. HDFS Sink を plugins にインストール
sudo -u kafka confluent-hub install --no-prompt \
confluentinc/kafka-connect-hdfs:latest \
--component-dir /opt/kafka/plugins
→最後のErrorが出ても無視してよい。
インストール結果確認:
find /opt/kafka/plugins -maxdepth 3 -type f -name '*.jar' | head
→/opt/kafka/plugins/confluentinc-kafka-connect-hdfs/lib/*.jarがあればよい。
- 本作業は基本オンラインでの作業が必要だが、それができない場所で実施する際は以下の手順で対応すること。
予め別のホストでConfluent-hub導入,プラグイン導入をすること。
cd /opt/kafka/plugins
sudo tar czf /tmp/confluentinc-kafka-connect-hdfs.tar.gz \
confluentinc-kafka-connect-hdfs
ここで作成されるconfluentinc-kafka-connect-hdfs.tar.gzをリポジトリサーバなどに保管し、以下のコマンドでプラグイン導入を実施する。
wget http://repo1/files/confluent/confluentinc-kafka-connect-hdfs.tar.gz
sudo tar xzf ./confluentinc-kafka-connect-hdfs.tar.gz \
-C /opt/kafka/plugins
5-3. Connect worker に plugin.path を設定
sudo grep -n '^plugin\.path' /opt/kafka/config/connect-distributed.properties || true
sudo sed -i '/^plugin\.path=/d' /opt/kafka/config/connect-distributed.properties
echo 'plugin.path=/opt/kafka/plugins' | sudo tee -a /opt/kafka/config/connect-distributed.properties
5-4. Kafka Connect を再起動
sudo systemctl restart kafka-connect 2>/dev/null || true
5-5. HDFS Sink が見えるか確認
curl -s http://localhost:8083/connector-plugins | jq . | grep -i hdfs -n || true
リストにio.confluent.connect.hdfs.HdfsSinkConnector が出ればよい。
6. guava.jar をプラグインに追加
6-1. OS の guava jar を使う
sudo apt update
sudo apt install -y libguava-java
インストールされる jar を確認
dpkg -L libguava-java | grep -E '/guava.*\.jar$'
以下のパスが出ることを確認
/usr/share/java/guava.jar
6-2. プラグインの lib にコピー
sudo cp -a /usr/share/java/guava*.jar /opt/kafka/plugins/confluentinc-kafka-connect-hdfs/lib/
sudo chown -R kafka:kafka /opt/kafka/plugins/confluentinc-kafka-connect-hdfs/lib
6-3. Connect を再起動
sudo systemctl restart kafka-connect 2>/dev/null || true
6-4. Connector作成
cat > hdfs-sink-syslog.json <<'EOF'
{
"name": "hdfs-sink-syslog",
"config": {
"connector.class": "io.confluent.connect.hdfs.HdfsSinkConnector",
"topics": "syslog",
"tasks.max": "1",
"hdfs.url": "hdfs://cluster1",
"hadoop.conf.dir": "/etc/hadoop/conf",
"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",
"topics.dir": "/data/kafka",
"flush.size": "500",
"rotate.interval.ms": "600000",
"rotate.schedule.interval.ms": "600000",
"schema.compatibility": "NONE"
}
}
EOF
curl -X POST http://localhost:8083/connectors \
-H "Content-Type: application/json" \
--data @hdfs-sink-syslog.json
curl -s http://localhost:8083/connectors/hdfs-sink-syslog/status | jq .
→stateがRUNNINGになること。
- 参考(本設定を書き換え再反映する場合)
cat > hdfs-sink-syslog.config.json <<'EOF'
{
"connector.class": "io.confluent.connect.hdfs.HdfsSinkConnector",
"topics": "syslog",
"tasks.max": "1",
"hdfs.url": "hdfs://cluster1",
"hadoop.conf.dir": "/etc/hadoop/conf",
"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",
"topics.dir": "/data/kafka",
"flush.size": "500",
"rotate.interval.ms": "600000",
"rotate.schedule.interval.ms": "600000",
"schema.compatibility": "NONE"
}
EOF
curl -X PUT http://localhost:8083/connectors/hdfs-sink-syslog/config -H "Content-Type: application/json" --data @hdfs-sink-syslog.config.json | jq .
curl -X POST http://localhost:8083/connectors/hdfs-sink-syslog/restart
curl -s http://localhost:8083/connectors/hdfs-sink-syslog/status | jq .
→stateがRUNNINGになること。
7. Zeppelin(Hive)DDL & 可視化
一度テーブルとVIEWを作成した後は、SELECTのみ実行すればよい。
各ホストで以下のコマンドを実行すると、数分後にHDFSへ保存され、Zeppelinから確認できる。
logger -p user.info "syslog forward test from $(hostname) $(date -Is)"
7-1. Hiveテーブル・VIEW作成
Kafka
Connectのvalue.converterにJsonConverterを使用しているため、HDFSに保存されたJSONをそのままget_json_objectで解析する。
%hive
CREATE DATABASE IF NOT EXISTS logs;
CREATE EXTERNAL TABLE IF NOT EXISTS logs.syslog (
line STRING
)
STORED AS TEXTFILE
LOCATION '/data/kafka/syslog';
CREATE OR REPLACE VIEW logs.v_syslog_parsed AS
WITH p AS (
SELECT
get_json_object(line, '$.src_host') AS host,
get_json_object(line, '$.ts') AS ts,
get_json_object(line, '$.message') AS msg,
get_json_object(line, '$.program') AS program
FROM logs.syslog
WHERE line IS NOT NULL
AND length(line) > 0
)
SELECT
host,
ts,
substr(ts, 1, 4) AS year,
substr(ts, 1, 7) AS month,
substr(ts, 1, 10) AS day,
substr(ts, 1, 13) AS hour,
substr(ts, 1, 16) AS minute,
substr(ts, 1, 19) AS sec,
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 p
WHERE host IS NOT NULL
AND ts IS NOT NULL
AND length(ts) >= 19;
7-2. WARNING以上のログ件数確認
%hive
set mapreduce.map.memory.mb=512;
set mapreduce.reduce.memory.mb=512;
set yarn.app.mapreduce.am.resource.mb=512;
set hive.mapred.supports.subdirectories=true;
set mapreduce.input.fileinputformat.input.dir.recursive=true;
SELECT
host,
hour,
COUNT(*) AS cnt
FROM logs.v_syslog_parsed
WHERE day >= date_format(current_date, 'yyyy-MM-dd')
AND severity <= 4
GROUP BY host, hour
ORDER BY hour, host;
TableからGraph(BarまたはPie)へ切り替え、当日毎時のWARNING以上のログ数が表示されることを確認する。
7-3. テストログ確認
%hive
set mapreduce.map.memory.mb=512;
set mapreduce.reduce.memory.mb=512;
set yarn.app.mapreduce.am.resource.mb=512;
set hive.mapred.supports.subdirectories=true;
set mapreduce.input.fileinputformat.input.dir.recursive=true;
SELECT
host,
ts,
program,
msg
FROM logs.v_syslog_parsed
WHERE lower(msg) LIKE '%syslog forward test from%'
ORDER BY ts DESC;
手順の最初にloggerで送信したログが表示されることを確認する。
確認実施後は以下のSQLを動かし、後始末を行う。
%hive
drop view logs.v_syslog_parsed;
drop table logs.syslog;
drop database logs;
運用上の補足
本記事のZeppelin VIEWは、構築直後の動作確認を目的としたものです。
実運用では以下の構成を推奨します。HDFS(JSON)→ Hive RAW → Hive Curated → Iceberg → Trino → Grafana
RAWテーブルではJSONを保持し、Hive Curatedで正規化・型変換を行った後、
Icebergを分析基盤として利用することで、大量データでも効率的に検索・集計できます。
8. kafka溜め込み確認
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe \
--group connect-hdfs-sink-syslog
→LAGの数が少ない、または減っていること。
なかなか減りにくい場合は、Kafkaのメモリを増強するかして確認する。
9. 今後の利用について
Zeppelinで作成したVIEWは構築直後の動作確認用途として利用します。
実運用では
HDFS(JSON)
↓
Hive RAW
↓
Hive Curated
↓
Iceberg
↓
Trino
↓
Grafana
の流れを推奨します。
RAWテーブルではJSON保持のみを行い、Hive Curatedで正規化したデータをIcebergへ取り込むことで、TrinoやGrafanaから高速に検索・分析できます。
詳細なCuratedテーブルの作成方法は次の記事を参照してください。