Kafka在CentOS生态中作为数据流通中枢,与EFK日志收集、HDFS存储、HBase、Prometheus+Grafana监控及SparkStreaming流处理系统协同,通过生产者-消费者模式构建实时数据管道,实现解耦、削峰填谷与高效集成。
Kafka在CentOS生态里扮演的角色远不只是个消息队列那么简单——它更像是整个数据流通的“中枢神经”。把Kafka和CentOS上的其他服务(比如日志收集、数据存储、NoSQL、监控系统)串起来,才能搭出真正意义上的实时数据管道。下面挑几种最常见的协同场景,逐一拆解具体怎么落地。

长期稳定更新的攒劲资源: >>>点此立即查看<<<
日志处理这块,EFK在CentOS上几乎是标配方案。Kafka夹在中间做缓冲,核心价值就是解耦——日志采集端不用管存储端是否扛得住,存储端也不用担心采集端突发流量。
nginx-logs,分区数设3、副本数设2;③ 修改Filebeat的配置文件(/etc/filebeat/filebeat.yml),把输出指向Kafka:output.kafka:
enabled: true
hosts: ["kafka1.centos:9092", "kafka2.centos:9092"]
topic: "nginx-logs"④ 启动服务:systemctl start filebeat kafka elasticsearch;⑤ 如果需要Logstash做中间处理,再加一条消费管道:bin/logstash -f /etc/logstash/conf.d/kafka-to-es.conf;⑥ 最后在Kibana里验证Nginx日志是否正常展示。把实时数据流灌进HDFS,是离线分析和历史存储的常见需求。Kafka到HDFS这条通路,通常靠Spark Streaming或Flume来搭桥。
order-data,分区数按数据量和消费并行度来设;③ 用Spark Streaming写一段消费程序(Java或Python),核心逻辑长这样:val sparkConf = new SparkConf().setAppName("KafkaToHDFS")
val ssc = new StreamingContext(sparkConf, Seconds(10))
val kafkaParams = Map[String, Object](
"bootstrap.servers" -> "kafka1.centos:9092,kafka2.centos:9092",
"key.deserializer" -> classOf[StringDeserializer],
"value.deserializer" -> classOf[StringDeserializer],
"group.id" -> "hdfs-writer",
"auto.offset.reset" -> "latest")
val topics = Array("order-data")
val stream = KafkaUtils.createDirectStream[String, String](ssc,
PreferConsistent, Subscribe[String, String](topics, kafkaParams))
stream.map(record => record.value()).saveAsTextFile("hdfs://namenode:8020/user/hadoop/order-data")
ssc.start()
ssc.awaitTermination()④ 提交Spark作业,然后在HDFS里确认数据文件是否生成。有些场景要求低延迟的随机读写,HBase正好对路。Kafka生产数据,HBase负责存储,中间还是靠消费者来牵线。
hbase-site.xml,把Zookeeper地址和HDFS路径填对:
hbase.rootdir
hdfs://namenode:8020/hbase
hbase.zookeeper.quorum
zookeeper1.centos:2181,zookeeper2.centos:2181
③ 写一个生产者程序,往Topic(比如user-data)里发数据;④ 消费者这边用Java实现,核心片段:Configuration config = HBaseConfiguration.create();
try (Connection connection = ConnectionFactory.createConnection(config);
Table table = connection.getTable(TableName.valueOf("user_table"))) {
KafkaConsumer consumer = new KafkaConsumer<>(kafkaProps);
consumer.subscribe(Arrays.asList("user-data"));
while (true) {
ConsumerRecords records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord record : records) {
Put put = new Put(Bytes.toBytes(record.key()));
put.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("info"), Bytes.toBytes(record.value()));
table.put(put);
}
}
} ⑤ 启动HBase和Kafka,跑通生产者和消费者,再查一下HBase表里有没有数据。Kafka跑在CentOS上,没有监控就等于裸奔。常用的监控组合是Kafka Exporter + Prometheus + Grafana,另外Kafka Manager和Burrow也各有用途。
./kafka_exporter --kafka.server=kafka1.centos:9092 --web.listen-address=:9308② 在Prometheus配置文件prometheus.yml里加上抓取任务:scrape_configs:
- job_name: 'kafka'
static_configs:
- targets: ['kafka1.centos:9308']③ Grafana里配好Prometheus数据源,导入一个Kafka仪表盘(比如社区ID 3955);④ 重点盯几个指标:吞吐量(kafka_server_brokertopicmetrics_messages_in_total)、消费滞后(kafka_consumer_fetch_manager_metrics_records_lag)、分区离线数(kafka_controller_kafkacontroller_offline_partitions_count)。实时统计、ETL之类的活儿,Spark Streaming是最常见的搭档。Kafka当数据源,Spark Streaming做计算,结果可以写回HDFS、数据库或者另一个Kafka Topic。
createDirectStream直接拉取Kafka数据,处理完再输出。spark-submit --class com.example.KafkaToHDFS --master yarn --deploy-mode cluster kafka-to-hdfs.jar;④ 去HDFS或者Spark UI上看一下结果是否如预期。以上这几个场景基本覆盖了Kafka在CentOS生态里最常见的搭档方式。关键在于理解每个组件在管道里扮演的角色,以及连接它们的套路——无非是生产者把数据送进Topic,消费者从Topic里拿出来,交给下游处理。把这个逻辑吃透了,剩下就是配置细节和调试的问题了。
侠游戏发布此文仅为了传递信息,不代表侠游戏网站认同其观点或证实其描述