在 Debian 上部署的 Kafka 对接外部系统,通常有三种主要方式:Kafka Connect、Producer/Consumer API/客户端、以及第三方中间件/流式计算框架。下面按常见场景给你一个清晰指南。
Kafka Connect 是官方提供的可扩展数据集成框架,适合数据库、消息队列、对象存储等。
如果你是用官方 Kafka 包:
sudo apt install default-jre
wget https://downloads.apache.org/kafka/3.x.x/kafka_2.13-x.x.x.tgz
tar -xzf kafka_*.tgz
cd kafka_*
启动:
bin/connect-distributed.sh config/connect-distributed.properties
{
"name": "mysql-source",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "mysql-host",
"database.user": "user",
"database.password": "pass",
"database.server.id": "184054",
"database.server.name": "dbserver",
"database.include.list": "test"
}
}
{
"name": "es-sink",
"config": {
"connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
"topics": "topic_name",
"connection.url": "http://es-host:9200"
}
}
适合自研系统或实时业务。
pip install kafka-python
from kafka import KafkaProducer
producer = KafkaProducer(bootstrap_servers='localhost:9092')
producer.send('test', b'hello')
from kafka import KafkaConsumer
consumer = KafkaConsumer('test', bootstrap_servers='localhost:9092')
for msg in consumer:
print(msg.value)
支持语言:
示例(Flink SQL):
CREATE TABLE kafka_source (
id INT,
name STRING
) WITH (
'connector' = 'kafka',
'topic' = 'test',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
);
| 外部系统 | 对接方式 |
|---|---|
| RabbitMQ | Kafka Connect / 自写桥接 |
| Redis | Consumer + Redis Client |
| S3 / MinIO | S3 Sink Connector |
| HTTP API | 自定义 Consumer |
| Logstash | Kafka 输入输出插件 |
sudo ufw allow 9092
listeners=PLAINTEXT://0.0.0.0:9092
advertised.listeners=PLAINTEXT://your-host-ip:9092
如果你能告诉我:
我可以给你更精确的配置示例。