如何用脚本批量生产Kafka消息?

wen 实用脚本 3

本文目录导读:

如何用脚本批量生产Kafka消息?

  1. 方法一:使用 Bash 脚本 + kafka-console-producer(最快速,适合简单文本)
  2. 方法二:使用 Python 脚本(最灵活,适合复杂数据、JSON、异步)
  3. 方法三:高级技巧——生成超大负载或特定性能测试
  4. 方法四:使用 kafkacat(命令行工具,轻量级)
  5. 总结与建议
  6. 调试小技巧

我们可以使用多种脚本工具来批量生产 Kafka 消息,最常见的方案是使用 Bash 脚本配合 Kafka 自带的命令行工具 kafka-console-producer,或者使用 Python 脚本结合 confluent-kafkakafka-python 库。

以下是几种最实用、最灵活的脚本方法:

使用 Bash 脚本 + kafka-console-producer(最快速,适合简单文本)

如果你只需要快速发送大量相同格式的文本消息,这是最简单的方法,它无需安装任何额外的库。

循环发送固定消息

#!/bin/bash
# 文件名: produce_fixed.sh
KAFKA_BROKER="localhost:9092"
TOPIC="my-topic"
MESSAGE_COUNT=1000
for i in $(seq 1 $MESSAGE_COUNT)
do
  # 发送消息到 Kafka
  echo "这是第 $i 条测试消息" | kafka-console-producer \
    --broker-list $KAFKA_BROKER \
    --topic $TOPIC \
    --property "parse.key=false" # 不发送 key
  # 可选:显示进度
  if (( $i % 100 == 0 )); then
    echo "已发送 $i 条消息"
  fi
done
echo "完成!共发送 $MESSAGE_COUNT 条消息"

从文件读取内容发送

#!/bin/bash
# 文件名: produce_from_file.sh
KAFKA_BROKER="localhost:9092"
TOPIC="my-topic"
# 假设我们有一个 large_messages.txt 文件,每行是一条消息
cat large_messages.txt | kafka-console-producer \
  --broker-list $KAFKA_BROKER \
  --topic $TOPIC \
  --property "parse.key=false"

使用 Key 发送(用于分区)

#!/bin/bash
KAFKA_BROKER="localhost:9092"
TOPIC="my-topic"
for i in {1..100}
do
  # 使用制表符分隔 key 和 value
  echo -e "user_$((i % 5))\t订单数据第 $i 条" | kafka-console-producer \
    --broker-list $KAFKA_BROKER \
    --topic $TOPIC \
    --property "parse.key=true" \
    --property "key.separator=\t"
done

使用 Python 脚本(最灵活,适合复杂数据、JSON、异步)

Python 是目前最推荐的方案,尤其适合需要批量生成 JSON、CSV 或进行错误处理的高级场景。

首先安装库:

pip install kafka-python  # 或 confluent-kafka

批量发送结构化 JSON 消息

# 文件名: produce_json.py
import json
import time
from kafka import KafkaProducer
import random
# 配置生产者
producer = KafkaProducer(
    bootstrap_servers=['localhost:9092'],
    value_serializer=lambda v: json.dumps(v).encode('utf-8') # 自动序列化为 JSON
)
TOPIC = "my-topic-user-events"
NUM_MESSAGES = 10000  # 发送 1 万条
# 批量生成数据
for i in range(NUM_MESSAGES):
    message = {
        "user_id": random.randint(1, 1000),
        "action": random.choice(["login", "click", "purchase", "logout"]),
        "timestamp": int(time.time() * 1000),
        "value": round(random.uniform(10.0, 1000.0), 2),
        "index": i
    }
    # 发送(异步)
    future = producer.send(TOPIC, value=message)
    # 可选:等待确认(会降低速度)
    # record_metadata = future.get(timeout=10)
    if i % 1000 == 0:
        print(f"已发送 {i} 条消息")
# 等待所有未完成的消息发送完毕
producer.flush()
print(f"完成!共发送 {NUM_MESSAGES} 条消息")

使用 confluent-kafka(性能更高,适合超大批量)

# 文件名: produce_confluent.py
from confluent_kafka import Producer
import json
import time
def delivery_report(err, msg):
    """回调函数,用于确认消息发送状态"""
    if err is not None:
        print(f'消息发送失败: {err}')
    # else:
    #     print(f'消息发送到 {msg.topic()} [{msg.partition()}]')
conf = {
    'bootstrap.servers': 'localhost:9092',
    'queue.buffering.max.messages': 100000,  # 缓冲区大小
    'batch.num.messages': 1000,              # 批量发送大小
    'compression.type': 'gzip'               # 压缩
}
producer = Producer(conf)
TOPIC = "my-topic"
# 批量发送 10 万条
for i in range(100000):
    data = {"id": i, "message": f"测试数据 {i}", "timestamp": time.time()}
    # 异步发送,带回调
    producer.produce(
        TOPIC, 
        key=str(i % 5),      # 按 key 分区
        value=json.dumps(data).encode('utf-8'),
        callback=delivery_report
    )
# 轮询处理回调
producer.poll(0)
print("正在等待所有消息发送完成...")
producer.flush()  # 等待所有消息发送完成
print("全部发送完成!")

高级技巧——生成超大负载或特定性能测试

生成指定大小的消息(负载测试)

# 假设你需要每条消息正好 1KB
MESSAGE_SIZE = 1024  # 1KB
# 用字母 'x' 填充到指定大小
payload = "x" * MESSAGE_SIZE
# 或者使用随机字符
import random
import string
payload = ''.join(random.choices(string.ascii_letters + string.digits, k=MESSAGE_SIZE))
# 然后用上面任意一种方法发送这个 payload

使用多线程/进程加速(极大量数据时)

import threading
from kafka import KafkaProducer
import json
def send_batch(producer, topic, start, end):
    for i in range(start, end):
        message = {"index": i, "data": f"message_{i}"}
        producer.send(topic, value=message)
    producer.flush()
# 创建多个线程
threads = []
num_threads = 10
total_messages = 1000000  # 100 万条
messages_per_thread = total_messages // num_threads
for t in range(num_threads):
    start = t * messages_per_thread
    end = start + messages_per_thread
    thread = threading.Thread(
        target=send_batch,
        args=(producer, "my-topic", start, end)
    )
    threads.append(thread)
    thread.start()
# 等待所有线程完成
for thread in threads:
    thread.join()
print("百万级消息发送完成!")

使用 kafkacat(命令行工具,轻量级)

如果你喜欢命令行,kafkacat 是一个强大的 C 语言编写的 Kafka 工具,性能很高。

安装: brew install kafkacat (Mac) 或 apt-get install kafkacat (Linux)

使用方式:

# 从标准输入发送 1000 条
for i in {1..1000}; do echo "message_$i"; done | kafkacat -P -b localhost:9092 -t my-topic
# 从文件发送
kafkacat -P -b localhost:9092 -t my-topic -l my_messages.txt
# 发送 JSON
echo '{"id":1,"value":"test"}' | kafkacat -P -b localhost:9092 -t my-topic

总结与建议

场景 推荐方案 理由
极快速简单文本 Bash + kafka-console-producer 无需安装,一行命令
复杂结构化数据 Python + kafka-pythonconfluent-kafka 灵活,支持 JSON/CSV
性能测试/超大批量 Python + confluent-kafka (异步+批量压缩) 吞吐量最高
命令行爱好者 kafkacat 轻量级,无需写脚本

调试小技巧

  1. 先验证 Consumer 能收到: 在另一个终端运行 kafka-console-consumer --bootstrap-server localhost:9092 --topic my-topic --from-beginning,观察是否能收到你发送的消息。
  2. 查看发送速度: 在脚本中加入时间统计,如 start_time = time.time(),结束时计算总耗时和每秒消息数。
  3. 错误处理: 生产环境中建议捕获 KafkaError 并重试失败的批次。

如果你有具体的需求(如特定的数据格式、目标吞吐量等),可以进一步描述,我可以帮你定制脚本!

抱歉,评论功能暂时关闭!