本文目录导读:

- 方法一:使用 Bash 脚本 + kafka-console-producer(最快速,适合简单文本)
- 方法二:使用 Python 脚本(最灵活,适合复杂数据、JSON、异步)
- 方法三:高级技巧——生成超大负载或特定性能测试
- 方法四:使用 kafkacat(命令行工具,轻量级)
- 总结与建议
- 调试小技巧
我们可以使用多种脚本工具来批量生产 Kafka 消息,最常见的方案是使用 Bash 脚本配合 Kafka 自带的命令行工具 kafka-console-producer,或者使用 Python 脚本结合 confluent-kafka 或 kafka-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-python 或 confluent-kafka |
灵活,支持 JSON/CSV |
| 性能测试/超大批量 | Python + confluent-kafka (异步+批量压缩) |
吞吐量最高 |
| 命令行爱好者 | kafkacat |
轻量级,无需写脚本 |
调试小技巧
- 先验证 Consumer 能收到: 在另一个终端运行
kafka-console-consumer --bootstrap-server localhost:9092 --topic my-topic --from-beginning,观察是否能收到你发送的消息。 - 查看发送速度: 在脚本中加入时间统计,如
start_time = time.time(),结束时计算总耗时和每秒消息数。 - 错误处理: 生产环境中建议捕获
KafkaError并重试失败的批次。
如果你有具体的需求(如特定的数据格式、目标吞吐量等),可以进一步描述,我可以帮你定制脚本!