1、新增功能探针联动处置、心跳在线检测
2、syslog-consumer模块拆分 syslog-consumer-rule模块实现日志数据消费、解析、泛化入库。
This commit is contained in:
+158
@@ -0,0 +1,158 @@
|
||||
package com.kafka;
|
||||
import com.common.entity.SyslogMessage;
|
||||
import com.common.entity.XdrHoneypot;
|
||||
import com.common.mapper.XdrHoneypotMapper;
|
||||
import com.common.util.MyBatisUtil;
|
||||
import com.common.util.SyslogParser;
|
||||
import com.influxdb.client.domain.WritePrecision;
|
||||
import com.influxdb.client.write.Point;
|
||||
import org.apache.ibatis.session.SqlSession;
|
||||
import org.apache.kafka.clients.consumer.*;
|
||||
import org.apache.kafka.common.serialization.StringDeserializer;
|
||||
import java.time.Duration;
|
||||
import java.time.format.DateTimeFormatter;
|
||||
import java.util.*;
|
||||
|
||||
import com.influx.InfluxDBClient;
|
||||
import com.common.util.JsonParser;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import com.common.util.JsonParser;
|
||||
import com.config.AppConfig;
|
||||
import com.Modules.NormalData.LogNormalProcessor;
|
||||
import java.time.LocalDate;
|
||||
import java.time.format.DateTimeFormatter;
|
||||
|
||||
public class kafkalogconsumer {
|
||||
|
||||
private static final Logger logger = LoggerFactory.getLogger(kafkalogconsumer.class);
|
||||
private static final DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
|
||||
private static final Random random = new Random();
|
||||
|
||||
public kafkalogconsumer() {
|
||||
Run();
|
||||
}
|
||||
|
||||
public static void main(String[] args) {
|
||||
Run();
|
||||
}
|
||||
public static void Run()
|
||||
{
|
||||
// 配置消费者属性
|
||||
//Properties props = new Properties();
|
||||
//props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "192.168.222.130:9092");
|
||||
// props.put(ConsumerConfig.GROUP_ID_CONFIG, "test-group-app");
|
||||
Properties props = new Properties();
|
||||
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, AppConfig.getBootstrapServers());
|
||||
props.put(ConsumerConfig.GROUP_ID_CONFIG, AppConfig.getGroupId());
|
||||
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
|
||||
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
|
||||
|
||||
// 可选配置
|
||||
//props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "none"); // 从最早的消息开始消费
|
||||
//props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // 从最早的消息开始消费
|
||||
//props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true"); // 自动提交偏移量
|
||||
//props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "1000"); // 自动提交间隔
|
||||
//props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 1000); // 设置单次拉取最大消息数[citation:6]
|
||||
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, AppConfig.getAutoOffsetReset()); // 从last开始消费
|
||||
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, AppConfig.getEnableAutoCommit()); // 自动提交偏移量
|
||||
props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG,AppConfig.getAutoCommitIntervalMS()); // 自动提交间隔
|
||||
|
||||
// 创建消费者实例
|
||||
Consumer<String, String> consumer = new KafkaConsumer<>(props);
|
||||
try {
|
||||
// 订阅主题
|
||||
consumer.subscribe(Collections.singletonList(AppConfig.getTopic()));
|
||||
|
||||
System.out.println("开始消费消息...");
|
||||
com.influx.InfluxDBClient influxClient = new InfluxDBClient();
|
||||
// 持续消费消息
|
||||
while (true) {
|
||||
// 拉取消息(等待最多100毫秒)
|
||||
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
|
||||
for (ConsumerRecord<String, String> record : records) {
|
||||
logger.info("收到syslogmessage:"+ record.value());
|
||||
System.out.printf(
|
||||
"收到消息: 主题=%s, 分区=%d, 偏移量=%d, 键=%s, 值=%s%n",
|
||||
record.topic(),
|
||||
record.partition(),
|
||||
record.offset(),
|
||||
record.key(),
|
||||
record.value()
|
||||
);
|
||||
String sysLogUUID =getSysLogUUID();
|
||||
String strDeviceInfo=SyslogParser.substringBeforeFirstChar(record.value(),']');
|
||||
Map<String,String> mapdev =SyslogParser.parseKeyValuePairs(strDeviceInfo);
|
||||
|
||||
// 初始化 InfluxDB 客户端
|
||||
Point point = Point.measurement("syslog_security")
|
||||
.addTag("deviceid", mapdev.get("device_id")) // 添加标签
|
||||
.addTag("uuid", sysLogUUID) //syslog uuid
|
||||
.addTag("topic", AppConfig.getTopic()) //kafka topic
|
||||
.addField("message", record.value()) // 添加字段
|
||||
.time(System.currentTimeMillis(), WritePrecision.MS) ;// 毫秒级时间戳
|
||||
influxClient.writePointBlocking(point);
|
||||
System.out.println("influxdb wirte syslog ,value:"+ record.key());
|
||||
// 日志信息插入pg XdrHoneypot 表
|
||||
//insertSingleRecord( record.value());
|
||||
System.out.println("insert postgres syslog ,value:"+ record.key());
|
||||
|
||||
//String syslogMessage= AppConfig.geRunEnvironment().equals("test")? record.value().substring(34) : record.value();
|
||||
String syslogMessage= record.value();
|
||||
//剔除测试环境本机syslog新增的头部信息
|
||||
LogNormalProcessor logNormalProcessor = new LogNormalProcessor(syslogMessage,sysLogUUID,AppConfig.getTopic());
|
||||
//LogNormalProcessor logNormalProcessor =new LogNormalProcessor(record.value());
|
||||
logNormalProcessor.init();
|
||||
}
|
||||
// 手动提交偏移量(如果禁用自动提交)
|
||||
consumer.commitSync();
|
||||
}
|
||||
} catch (Exception e) {
|
||||
e.printStackTrace();
|
||||
} finally {
|
||||
// 关闭消费者
|
||||
consumer.close();
|
||||
}
|
||||
}
|
||||
/**
|
||||
* 获取日志信息UUID,格式: yyyyMMddxxxxxxxx
|
||||
* @return
|
||||
*/
|
||||
private static String getSysLogUUID ()
|
||||
{
|
||||
// 获取当前日期
|
||||
LocalDate currentDate = LocalDate.now();
|
||||
// 定义格式 (yyyyMMdd)
|
||||
DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyyMMdd");
|
||||
// 格式化日期
|
||||
String formattedDate = currentDate.format(formatter);
|
||||
return formattedDate +"-"+ UUID.randomUUID() ;
|
||||
}
|
||||
|
||||
/**
|
||||
* 单条记录插入演示
|
||||
*/
|
||||
private static void insertSingleRecord(String strlog) {
|
||||
logger.info("=== 单条记录插入演示 ===");
|
||||
|
||||
try (SqlSession sqlSession = MyBatisUtil.getSqlSession()) {
|
||||
XdrHoneypotMapper mapper = sqlSession.getMapper(XdrHoneypotMapper.class);
|
||||
|
||||
// 创建测试数据
|
||||
//XdrHoneypot record = createTestXdrHoneypot(1);
|
||||
XdrHoneypot record = JsonParser.parseLogMessageToXdrHoneypot( strlog);
|
||||
// 插入记录
|
||||
int result = mapper.insert(record);
|
||||
sqlSession.commit();
|
||||
|
||||
if (result > 0) {
|
||||
logger.info("单条记录插入成功,ID: {}", record.getId());
|
||||
logger.info("插入的数据: {}", record);
|
||||
} else {
|
||||
logger.error("单条记录插入失败");
|
||||
}
|
||||
} catch (Exception e) {
|
||||
logger.error("单条记录插入出错", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user