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?

Kafka MirrorMaker2によるHadoopクラスタ間レプリケーションやってみた。

0
Last updated at Posted at 2026-01-03

はじめに

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:過去のデータ&最新のデータが参照できること。

0
0
3

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?