Lua将Nginx请求数据写入Kafka——埋点日志解决方案

缘起

有一个埋点收集系统,架构是Nginx+Flume。web,小程序,App等客户端将数据报送至Nginx,Nginx将请求写入本地文件,然后Flume读取日志文件的数据,将日志写入Kafka。这个架构本来没什么问题,但是部署在K8s容器就有问题了,当前一个Nginx后面是3个Flume,Nginx根据渠道将日志写入web.log,mp.log,app.log,3个log文件各对应一个Flume将数据写入Kafka,遇到的问题首先是health check问题,K8s一个容器只能提供1个health check地址(最佳实战),因为Nginx后面有3个Flume,所以无法感知这3个健康状态,最多感知一个。第二个动态扩容问题,比如某一段时间web报送数据较多,但其它两个渠道较少,扩容成两个pod会浪费资源,第三个是优雅退出问题,由于Nginx速度较快,写入文件也比较快,Flume处理的较慢,导致如果我想把这个容器关闭的话,不知道Flume有没有把所有的日志都写入Kafka。 优化思路:Flume其实是有多个进口和多个出口的,对于客户的业务来讲,入口只有Nginx,出口只有Kafka,于是决定将Flume去掉,使用Nginx+Lua脚本形式直接将Nginx的日志写入Kafka,同时写一份文件以备出现问题补数据使用,也作为Bug定位追踪使用。

环境准备

  1. openresty-1.15.8.2 下载地址:http://openresty.org/cn/download.html
  2. lua-resty-kafka 是lua版本的kafka驱动,内置producer 下载地址:https://github.com/doujiang24/lua-resty-kafka
  3. zlib的lua库 解gzip使用的,如果你不用gzip 可以不安装到脚本 下载地址:https://github.com/madler/zlib
  4. lua-zlib lua调用gzip使用的库 下载地址:https://github.com/brimworks/lua-zlib
  5. Nginx和上面组件编译用的依赖:(我用的debian10.11)
1apt-get update -y && apt-get install --fix-missing zlib1g zlib1g-dev libpcre3-dev libssl-dev perl make build-essential curl cmake -y

步骤

  1. 解压openresty并编译安装
1cp openresty-1.15.8.2.tar.gz /opt/app/source/ 2tar -zxvf openresty-1.15.8.2.tar.gz 3 4cd openresty-1.15.8.2 && ./configure --prefix=/opt/app/openresty && make && make install
  1. 将kafka驱动放入lualib 就是将lua-resty-kafka-0.10.zip解压开,把lua-resty-kafka-0.10/lib下的resty文件夹直接拷贝进/opt/app/openresty/site/lualib/下
1unzip lua-resty-kafka-0.10.zip 2cd lua-resty-kafka-0.10/lib 3cp -r resty /opt/app/openresty/site/lualib/
  1. 安装lua写的zlib库和lua-zlib库
1tar -zxvf zlib-master.tar.gz 2cp -a zlib-master/* /opt/app/openresty/site/lualib/ 3 4tar -zxvf lua-zlib-master.tgz 5cd lua-zlib-master \ 6 && cmake -DLUA_INCLUDE_DIR=/opt/app/openresty/luajit/include/luajit-2.1 -DLUA_LIBRARIES=/opt/app/openresty/luajit/lib -DUSE_LUAJIT=ON -DUSE_LUA=OFF \ 7 && make \ 8 && cp zlib.so /opt/app/openresty/lualib/zlib.so
  1. 编写kafkaconfig.lua并放入/opt/app/openresty/site/lualib/ 这一步是为了配置kafka的ip,端口和topic。 内容是:
1kafka_broker_list={ 2 {host="192.168.1.1",port=9092}, 3 {host="192.168.1.2",port=9092}, 4 {host="192.168.1.3",port=9092} 5} 6 7kafka_topic_mp="mp_topic" 8kafka_topic_app="app_topic" 9kafka_topic_web="web_topic"
  1. 编写nginx配置文件nginx.conf.内容如下:
1worker_processes auto; 2error_log /opt/app/logs/error.log; 3events { 4 worker_connections 10240; 5} 6 7http { 8 gzip on; 9 gzip_types application/javascript text/plain text/xml text/css application/x-javascript application/xml text/javascript application/x-httpd-php image/jpeg image/gif image/png; 10 gzip_vary on; 11 gzip_comp_level 9; 12 13 server { 14 listen 80; 15 include /opt/app/openresty/nginx/conf/web.conf; 16 } 17 18 } 19}

web.conf

1location ^~ /health { 2 default_type text/html; 3 return 200 'ok'; 4} 5# app 批量 6location ^~ /app { 7 if ($request_method != POST) { 8 return 405; 9 } 10 default_type 'text/plain'; 11 expires off; 12 add_header Last-Modified ''; 13 add_header Cache-Control 'no-cache'; 14 add_header Pragma "no-cache"; 15 empty_gif; 16 echo_read_request_body; 17 access_by_lua_file /opt/app/openresty/nginx/conf/app.lua; 18 access_log /opt/app/logs/app.log; 19}

app.lua:

1ngx.req.read_body() 2local body = ngx.req.get_body_data() 3 4#引入kafka的生产者 5local producer = require "resty.kafka.producer" 6#引入上面写的kafka配置 7local kafka_config = require "kafkaconfig" 8 9# 因为引入了kafkaconfig 里面的变量直接访问就好 10# local p = producer:new(kafka_broker_list) 11# --------> 2023-03-16 分界线 start<------------ 12# 请使用以下方式初始化producer 相比于上面的初始化,下面添加了broker刷新时间 13# 可以在网络抖动获取不到boker list,防止后续无法继续使用的问题 14# 详细见我另一篇文章 https://blog.csdn.net/codeblf2/article/details/129505283 15local p = producer:new(kafka_broker_list, {producer_type = "async",refresh_interval=10000}) 16# --------> 2023-03-16 分界线 end <------------ 17 18# 中间为nil的参数是kafka将本条消息写入哪个分区所使用的key,为nil代表轮询写入 19local offset, err = p:send(kafka_topic_app, nil, body) 20if not offset then 21 ngx.say("send err:", err) 22 return 23end

以上示例只展示了app的,web和mp的类似即可。

点赞
收藏

评论区

加载中...

相关推荐

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

Nginx 502 Bad Gateway 的错误的解决方案

我用的是nginx反向代理Apache,直接用Apache不会有任何问题,加上nginx就会有部分ajax请求502的错误,下面是我收集到的解决方案。一、fastcgi缓冲区设置过小 出现错误,首先要查找nginx的日志文件,目录为/var/log/nginx,在日志中发现了如下错误 2013/01/1713:33:47\err

ELK之八

一、logstash结合kafka收集系统日志和nginx日志架构图:!(https://oscimg.oschina.net/oscnet/2d28dece38ea896fdb974165c799ff8130a.png)环境准备:A主机:kibana、e

Flume使用Kafka Sink导致CPU过高的问题

在日志收集服务器上使用Flume(1.6)的KafkaSink将日志数据发送至Kafka,在FlumeAgent启动之后,发现每个Agent的CPU使用率都非常高,而我们需要在每台机器上启动多个FlumeAgent来收集不同类型的日志,如果每个Agent都这样,那肯定会把机器的CPU吃满了,刚开始使用jstack定位到是org.apache.flume

JVM 字节码指令表

字节码助记符指令含义0x00nop什么都不做0x01aconst\_null将null推送至栈顶0x02iconst\_m1将int型1推送至栈顶0x03iconst\_0将int型0推送至栈顶0x04iconst\_1将int型1推送至栈顶0x05ic

埋点日志最终解决方案——Golang+Gin+Sarama VS Java+SpringWebFlux+ReactorKafka

埋点日志最终解决方案——GolangGinSaramaVSJavaSpringWebFluxReactorKafka之前我就写过几篇OpenRestyluakafkaclient将埋点数据写入Kafka的文章,如下:以上一步一个坑,有些是自己能力