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.common.util.SpringContextUtil; import com.config.AppProperties; import com.config.KafkaConsumerProperties; 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() { KafkaConsumerProperties kafkaProps = SpringContextUtil.getBean(KafkaConsumerProperties.class); LogNormalProcessor logNormalProcessor = SpringContextUtil.getBean(LogNormalProcessor.class); Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaProps.getBootstrapServers()); props.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaProps.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, kafkaProps.getAutoOffsetReset()); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, kafkaProps.isEnableAutoCommit()); props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, kafkaProps.getAutoCommitInterval()); // 创建消费者实例 Consumer consumer = new KafkaConsumer<>(props); try { // 订阅主题 consumer.subscribe(Collections.singletonList(kafkaProps.getTopic())); System.out.println("开始消费消息..."); com.influx.InfluxDBClient influxClient = SpringContextUtil.getBean(com.influx.InfluxDBClient.class); // 持续消费消息 while (true) { // 拉取消息(等待最多100毫秒) ConsumerRecords records = consumer.poll(Duration.ofMillis(100)); for (ConsumerRecord 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 mapdev =SyslogParser.parseKeyValuePairs(strDeviceInfo); // 初始化 InfluxDB 客户端 Point point = Point.measurement("syslog_security") .addTag("deviceid", mapdev.get("device_id")) // 添加标签 .addTag("uuid", sysLogUUID) //syslog uuid .addTag("topic", kafkaProps.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= record.value(); //使用注入的 Spring Bean 进行标准化处理 logNormalProcessor.process(syslogMessage, sysLogUUID, null); } // 手动提交偏移量(如果禁用自动提交) 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); } } }