159 lines
7.3 KiB
Java
159 lines
7.3 KiB
Java
|
|
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);
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
}
|