はじめに
Kafka MirrorMaker2によるHadoopクラスタ間レプリケーションをやってみました。
全体像
前提
以下の手順を参考にそれぞれのクラスタを構築すること。
src クラスタ
以下のページを参考に構築していること。
fluentd → Kafka(syslog topic)
Kafka broker: kafka1,kafka2,kafka3
dst クラスタ
以下のページを参考に構築していること。
Kafka + HDFS
Kafka broker: hdptest
MirrorMaker 専用ホスト(新規)
ホスト名:mm2-1
src / dst 両方の broker に疎通可能
1.dst クラスタの syslog 収集停止・データクリア
1-1. dst側 fluentd停止
sudo systemctl stop fluentd
確認:
systemctl status fluentd
1-2. dst側 HDFS Sink 停止
sudo systemctl stop kafka-connect
確認:
systemctl status kafka-connect
1-3. dst Kafka の syslog topic 削除
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server hdptest:9092 \
--delete \
--topic syslog
確認:
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server hdptest:9092 \
--list | grep syslog || echo "syslog deleted"
1-4. dst HDFS データ削除
sudo -u hadoop hdfs dfs -rm -r -skipTrash /data/kafka/syslog
sudo -u hadoop hdfs dfs -rm -r -skipTrash /data/kafka/+tmp/syslog
確認:
hdfs dfs -ls /data/kafka
hdfs dfs -ls /data/kafka/+tmp
1-5. dst側 HDFS Sink 開始
sudo systemctl start kafka-connect
確認:
systemctl status kafka-connect
hdfs dfs -ls -R -h /data/kafka/syslog | head -n 100
→HDFSに新たに溜め込まれるログが出ていないこと。
2.MirrorMaker ホスト構築
2-1. 必要パッケージ
sudo apt update
sudo apt install -y openjdk-17-jre-headless curl tar jq
sudo useradd -r -m -s /usr/sbin/nologin mm2 || true
2-2. Kafka 配布物配置(MM2-1 実行用)
cd /tmp
curl -LO 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 mm2:mm2 /opt/kafka
2-3. MM2-1 設定ファイル
sudo mkdir -p /etc/kafka-mm2
sudo chown -R mm2:mm2 /etc/kafka-mm2
sudo tee /etc/kafka-mm2/mm2.properties <<'EOF'
clusters = src, dst
# 複製元:3ブローカー
src.bootstrap.servers = kafka1:9092,kafka2:9092,kafka3:9092
# 複製先:1ブローカー
dst.bootstrap.servers = hdptest:9092
src->dst.enabled = true
src->dst.topics = ^(syslog)$
# 対象トピックが複数ある場合
# src->dst.topics = ^(syslog|authlog)$
src->dst.groups = .*
# 複製先でも元と同じトピック名を使用
replication.policy.class = org.apache.kafka.connect.mirror.IdentityReplicationPolicy
# 保存済みoffsetが存在しない場合は先頭から取得
#consumer.auto.offset.reset = earliest
# 保存済みoffsetが存在しない場合は末尾から取得
consumer.auto.offset.reset=latest
# --------------------------------------------------
# 複製先トピック
# dstは1台構成なので1
# --------------------------------------------------
replication.factor = 1
# --------------------------------------------------
# Kafka Connect内部トピック
# dstは1台構成なので1
# --------------------------------------------------
config.storage.replication.factor = 1
offset.storage.replication.factor = 1
status.storage.replication.factor = 1
# --------------------------------------------------
# MirrorMaker 2内部トピック
# checkpoint、heartbeatはdst側
# --------------------------------------------------
checkpoints.topic.replication.factor = 1
heartbeats.topic.replication.factor = 1
# offset-syncsはデフォルトではsrc側
# srcは3台構成なので3
offset-syncs.topic.location = source
offset-syncs.topic.replication.factor = 3
# REST API
rest.port = 18083
# トピック・Consumer Groupの再検出間隔
refresh.topics.interval.seconds = 30
refresh.groups.interval.seconds = 30
EOF
2-4. systemd 登録
sudo tee /etc/systemd/system/kafka-mm2.service <<'EOF'
[Unit]
Description=Kafka MirrorMaker 2
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
User=mm2
Group=mm2
WorkingDirectory=/opt/kafka
Environment="KAFKA_HEAP_OPTS=-Xms512m -Xmx512m"
ExecStart=/opt/kafka/bin/connect-mirror-maker.sh /etc/kafka-mm2/mm2.properties
Restart=on-failure
RestartSec=5
TimeoutStopSec=120
LimitNOFILE=100000
[Install]
WantedBy=multi-user.target
EOF
sudo systemctl daemon-reload
sudo systemctl enable --now kafka-mm2
3.MirrorMaker による src → dst レプリケーション
3-1. dst Kafka に topic ができるか確認
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server hdptest:9092 \
--list | grep syslog
3-2. dst Kafka にデータが来ているか
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server hdptest:9092 \
--topic syslog \
--from-beginning \
--max-messages 5
3-3. dst HDFSにデータが来ているか
hdfs dfs -ls -R -h /data/kafka/syslog | head -n 100
4. 追いつき確認 → src fluentd の送付先変更
SRC→DSTに切り替えするつもり無い(DSTバックアップ利用)の場合、実施しなくてもよい。
4-1. offset 突き合わせ
src/kctrl1
/opt/kafka/bin/kafka-get-offsets.sh \
--bootstrap-server kafka1:9092 \
--topic syslog
dst/hdptest
/opt/kafka/bin/kafka-get-offsets.sh \
--bootstrap-server hdptest:9092 \
--topic syslog
partition ごとに数字一致するまで何回かコマンド入力する。
4-2. src fluentd 停止
src/log1
sudo systemctl stop fluentd
30秒待つ。
4-3. 最終差分確認(もう一度 offset)
src/kctrl1
/opt/kafka/bin/kafka-get-offsets.sh \
--bootstrap-server kafka1:9092 \
--topic syslog
dst/hdptest
/opt/kafka/bin/kafka-get-offsets.sh \
--bootstrap-server hdptest:9092 \
--topic syslog
完全一致 = OK
4-4. src fluentd の送付先変更(dst Kafka へ)
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 hdptest:9092
default_topic syslog
<format>
@type json
</format>
required_acks 1
compression_codec gzip
</match>
EOF
反映:
sudo systemctl start fluentd
4-5. dst Kafka に直接流れているか確認
hdptestで実施
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server hdptest:9092 \
--topic syslog \
--max-messages 5
4-6. dst HDFS に流れているか確認
hdptestで実施
hdfs dfs -ls -R -h /data/kafka/syslog | head -n 100
5. MirrorMaker停止
5-1. MirrorMaker 停止
sudo systemctl stop kafka-mm2
sudo systemctl disable kafka-mm2
5-2. サービス・ファイル削除
sudo rm -f /etc/systemd/system/kafka-mm2.service
sudo rm -rf /etc/kafka-mm2
sudo rm -rf /opt/kafka
sudo userdel mm2
sudo systemctl daemon-reload
5-3. 最終確認
systemctl status kafka-mm2 || echo "MM2 removed"
上記作業後、本マシンを削除してもよい。
6. Zeppelin参照確認
src/dstそれぞれのZeppelinで以下のようになっていることを確認
-
src:過去のデータは保持されているが、最新データはないこと。
-
dst:過去のデータ&最新のデータが参照できること。