kafka最新版本(filebeat利用kafka进行日志实时传输)

本文目录
filebeat利用kafka进行日志实时传输
选择安装目录:例如安装在/usr/local/或者/opt/下都可以。
创建一个软链接:
filebeat的配置很简单,只需要指定input和output就可以了。
由于kafka server高低版本的客户端API区别较大,因此推荐同时使用高版本的filebeat和kafka server。 注意 高版本的filebeat配置使用低版本的kafka server会出现kafka server接受不到消息的情况。这里我使用的kafka server版本是:2.12-0.11.0.3,可参考 快速搭建kafka
filebeat安装目录下 filebeat.yml 文件:
配置Filebeat inputs:
上面 /opt/test/*.log 是我要传输的日志,根据实际情况改成你自己的值。
配置Filebeat outputs:
"111.11.1.243:9092" 是我的单机kafka broker,如果你是kafka集群,请用 , 分隔。 test 是kafka topic,请改成你自己实际情况的值。另外以下这段需要删除:
因为我并没有用到Elasticsearch,所以有多个输出在启动filebeat时会报错。到这里filebeat配置kafka就完成了,是不是很简单,让我们启动它测试一下。
启动,进入filebeat的安装目录:
查看是否启动:
很好,已经启动了。如果没有启动,请查看启动日志文件nohup.out。
停止:
随机生成日志脚本:
执行这段python脚本,开启一个kafka消费者如果成功消费日志消息:
哈哈,大功告成。 注 上面这段脚本要适时手动停止,因为它是个死循环,如果忘记手动停止那么就杯具了,我就是这样把机器写宕机的。
Flink kafka kerberos的配置
Flink消费集成kerberos认证的kafka集群时,需要做一些配置才可以正常执行。
Flink版本:1.8;kafka版本:2.0.1;Flink模式:Standalone
//指示是否从 Kerberos ticket 缓存中读取
security.kerberos.login.use-ticket-cache: false1
//Kerberos 密钥表文件的绝对路径
security.kerberos.login.keytab: /data/home/keytab/flink.keytab
//认证主体名称
security.kerberos.login.principal: flink@data.com
//Kerberos登陆contexts
security.kerberos.login.contexts: Client,KafkaClient
val properties: Properties =new Properties()
properties.setProperty("bootstrap.servers","broker:9092")
properties.setProperty("group.id","testKafka")
properties.setProperty("security.protocol","SASL_PLAINTEXT")
properties.setProperty("sasl.mechanism","GSSAPI")
properties.setProperty("sasl.kerberos.service.name","kafka")
consumer =new FlinkKafkaConsumer("flink",new SimpleStringSchema(), properties)
参数说明 :security.protocol
运行参数可以配置为PLAINTEXT(可不配置)/SASL_PLAINTEXT/SSL/SASL_SSL四种协议,分别对应Fusion Insight Kafka集群的21005/21007/21008/21009端口。 如果配置了SASL,则必须配置sasl.kerberos.service.name为kafka,并在conf/flink-conf.yaml中配置security.kerberos.login相关配置项。如果配置了SSL,则必须配置ssl.truststore.location和ssl.truststore.password,前者表示truststore的位置,后者表示truststore密码。
Kafka压缩
首先说明一点kafka的压缩和kafka的compact是不同的,compact就是相同的key只保留一条,是数据清理方式的一种和kafka的定期删除策略是并列的;而kafka的压缩是指数据不删除只是采用压缩算法进行压缩。
kafka从0.7版本就开始支持压缩功能:
1)kafka的发送端将消息按照批量(如果批量设置一条或者很小,可能有相反的效果)的方式进行压缩。
2)服务器端直接将压缩消息保存(特别注意,如果kafka的版本不同,那么就存在broker需要先解压缩再压缩的问题,导致消耗资源过多)。
3)消费端自动解压缩,测试了下,发送端无论采用什么压缩模式,消费端无论设置什么解压模式,都可以自动完成解压缩功能。
4)压缩消息可以和非压缩消息混存,也就是说如果你kafka里面先保存的是非压缩消息,后面改成压缩,不用担心,kafka消费端自动支持。
测试的kafka版本:kafka_2.12-1.1.1
测试的kafka客户端版本:0.10.2.1
测试数据的条数:20000
kafka支持三种压缩算法,lz4、snappy、gzip,
通过上面数据来看,gzip的压缩效果最好,但是生成耗时更长,snappy和lz4的数据差不多,更倾向于lz4,具huxi大神的书上所说kafka里面对snappy做了硬编码,所以性能上最好的是lz4,推荐使用此压缩算法。
压缩率对比:
性能对比图:
很简单:
消费端无论设置什么压缩模式,都可以正确的解压kafka的消息,也就是说消费端可以不设置解压缩,
不过可能性能有所下降。
Kafka的消息格式及offset是如何设置的
Kafka的offset是如何设置的?
答:是生产者设置的,生产者在发送消息的时候,为每条消息生成一个唯一的offset。
Kafka消息的格式?
答:
Kafka最新版本的消息集叫做RecordBatch,而不是先前的MessageSet。RecordBatch内部存储了一条或多条消息。
RecordBatch的结构包含以下部分:
first offset,起始位移,占位8B
length,消息总长度,占位4B
partition leader epoch,分区leader纪元,可以看做分区leader的版本号或者更新次数,占位4B。
magic,消息格式的版本号,对于V2版本而言,magic的值为2。
attributes,消息属性,占位2B,低三位表示压缩格式,第4位表示时间戳类型,第五位表示当前RecordBatch是否处于事务中第6位表示是否控制消息。
last offset delta,占位4B,RecordBatch中最后一个Record的offset与first offset的差值,主要被broker用来确保RecordBatch中Record组装的正确性。
first timestamp,占位8B,RecordBatch中第一条Record的时间戳。
max timestamp,占位8B,RecordBatch中最大的时间戳,一般情况下是最后一个Record的时间戳。和last offset delta功能一样,主要被broker用来确保RecordBatch中Record组装的正确性。
producer id,即PID,占位8B,用来支持幂等和事务。
producer epoch,占位2B,用来支持幂等和事务。
first sequence,占位4B,用来支持幂等和事务。
records count,占位4B,RecordBatch中Record的总数。
records,存放消息的容器。
Records的数据结构又是什么样的呢?Record包含以下属性:
length,消息总长度。
attributes,目前已弃用,但是还是会占用1B的空间,以备未来的格式扩展。
timestamp delta,时间戳增量,通常一个timestamp占用8B,这里时间戳增量保存的是当前时间戳与RecordBatch中first timestamp的差值,这样可以节省占用空间。
offset delta,位移增量,这个是当前消息的位移与RecordBatch中first offset的差值,这样可以节省占用空间。
key length,消息的key的长度。
key,消息key。
value length,消息的value的长度。
value,消息的值。
headers count,headers的总数。
headers,这个字段是用来支持应用级别的扩展,而不需要将一些应用级别的属性嵌入到消息体中。

更多文章:
数据库管理系统和数据库系统分别侧重(数据库,数据库管理系统,数据库系统,这三个分别是什么意思并举个实例)
2026年9月7日 17:00
springmvc的依赖(springMVC的注入方式有哪几种,这与springMVC依赖)
2026年9月7日 14:00
display flex 自动换行(overflow-y:hidden;overflow-x:auto;无效解决方法)
2026年9月7日 11:00
timestamp without time zone(Postgresql中to_date()函数使用问题)
2026年9月7日 09:40







