Skip to content

Commit af3d2b4

Browse files
Merge pull request #232 from shiyindaxiaojie/feature
Feature
2 parents 69b583c + 16a68ec commit af3d2b4

12 files changed

Lines changed: 153 additions & 32 deletions

File tree

eden-components/eden-solutions/eden-common-excel/src/main/java/org/ylzl/eden/common/excel/integration/easypoi/package-info.java

Lines changed: 0 additions & 1 deletion
This file was deleted.
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
package org.ylzl.eden.common.excel.integration.fastexcel;

eden-components/eden-solutions/eden-common-mq/src/main/java/org/ylzl/eden/common/mq/MessageQueueConsumer.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,8 @@ public interface MessageQueueConsumer {
3232
/**
3333
* 消费消息
3434
*
35-
* @param messages
35+
* @param messages 消息报文
36+
* @param ack 消息确认
3637
*/
3738
void consume(List<Message> messages, Acknowledgement ack);
3839
}

eden-components/eden-solutions/eden-common-mq/src/main/java/org/ylzl/eden/common/mq/integration/kafka/KafkaProvider.java

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -41,9 +41,9 @@
4141
@Slf4j
4242
public class KafkaProvider implements MessageQueueProvider {
4343

44-
private static final String KAFKA_PROVIDER_SEND_INTERRUPTED = "KafkaProvider send interrupted: {}";
44+
private static final String KAFKA_PROVIDER_SEND_INTERRUPTED = "KafkaProvider send interrupted";
4545

46-
private static final String KAFKA_PROVIDER_CONSUME_ERROR = "KafkaProvider send error: {}";
46+
private static final String KAFKA_PROVIDER_SEND_ERROR = "KafkaProvider send error";
4747

4848
private final KafkaTemplate<String, String> kafkaTemplate;
4949

@@ -70,11 +70,11 @@ public MessageSendResult syncSend(Message message) {
7070
SendResult<String, String> sendResult = future.get();
7171
return transfer(sendResult);
7272
} catch (InterruptedException e) {
73-
log.error(KAFKA_PROVIDER_SEND_INTERRUPTED, e.getMessage(), e);
73+
log.error(KAFKA_PROVIDER_SEND_INTERRUPTED, e);
7474
Thread.currentThread().interrupt();
7575
throw new MessageSendException(e.getMessage());
7676
} catch (Exception e) {
77-
log.error(KAFKA_PROVIDER_CONSUME_ERROR, e.getMessage(), e);
77+
log.error(KAFKA_PROVIDER_SEND_ERROR, e);
7878
throw new MessageSendException(e.getMessage());
7979
}
8080
}
@@ -102,7 +102,7 @@ public void onFailure(Throwable e) {
102102
}
103103
});
104104
} catch (Exception e) {
105-
log.error(KAFKA_PROVIDER_CONSUME_ERROR, e.getMessage(), e);
105+
log.error(KAFKA_PROVIDER_SEND_ERROR, e);
106106
throw new MessageSendException(e.getMessage());
107107
}
108108
}

eden-components/eden-solutions/eden-common-mq/src/main/java/org/ylzl/eden/common/mq/model/Message.java

Lines changed: 7 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -36,39 +36,25 @@
3636
@Data
3737
public class Message implements Serializable {
3838

39-
/**
40-
* 命名空间
41-
*/
39+
/** 命名空间 */
4240
private String namespace;
4341

44-
/**
45-
* 主题
46-
*/
42+
/** 主题 */
4743
private String topic;
4844

49-
/**
50-
* 分区/队列
51-
*/
45+
/** 分区/队列 */
5246
private Integer partition;
5347

54-
/**
55-
* 分区键
56-
*/
48+
/** 分区键 */
5749
private String key;
5850

59-
/**
60-
* 标签过滤
61-
*/
51+
/** 标签过滤 */
6252
private String tags;
6353

64-
/**
65-
* 消息体
66-
*/
54+
/** 消息体 */
6755
private String body;
6856

69-
/**
70-
* 延时等级
71-
*/
57+
/** 延时等级 */
7258
@Builder.Default
7359
private Integer delayTimeLevel = 0;
7460
}

eden-components/eden-solutions/eden-event-auditor/pom.xml

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,20 @@
6161
<optional>true</optional>
6262
</dependency>
6363

64+
<!-- Kafka -->
65+
<dependency>
66+
<groupId>org.springframework.kafka</groupId>
67+
<artifactId>spring-kafka</artifactId>
68+
<optional>true</optional>
69+
</dependency>
70+
71+
<!-- RocketMQ -->
72+
<dependency>
73+
<groupId>org.apache.rocketmq</groupId>
74+
<artifactId>rocketmq-spring-boot-starter</artifactId>
75+
<optional>true</optional>
76+
</dependency>
77+
6478
<!-- 测试组件 -->
6579
<dependency>
6680
<groupId>org.spockframework</groupId>

eden-components/eden-solutions/eden-event-auditor/src/main/java/org/ylzl/eden/event/auditor/aop/EventAuditorInterceptor.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -130,7 +130,7 @@ public void setAsyncTaskExecutor(AsyncTaskExecutor asyncTaskExecutor) {
130130
* @param events 审计事件列表
131131
*/
132132
private void send(List<AuditingEvent> events) {
133-
String senderType = eventAuditorConfig.getSender().getSenderType();
133+
String senderType = eventAuditorConfig.getSender().getType();
134134
EventSenderBuilder eventSenderBuilder = ExtensionLoader.getExtensionLoader(EventSenderBuilder.class).getExtension(senderType);
135135
eventSenderBuilder.setEventAuditorConfig(eventAuditorConfig);
136136
EventSender eventSender = eventSenderBuilder.build();

eden-components/eden-solutions/eden-event-auditor/src/main/java/org/ylzl/eden/event/auditor/config/EventAuditorConfig.java

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,7 @@ public class EventAuditorConfig {
2626
@Getter
2727
public static class Sender {
2828

29-
private String senderType = "logging";
29+
private String type = "logging";
3030

3131
private boolean async = true;
3232

@@ -51,6 +51,7 @@ public static class Logging {
5151
@Getter
5252
public static class Kafka {
5353

54+
private String topic;
5455
}
5556

5657
@EqualsAndHashCode
@@ -59,6 +60,13 @@ public static class Kafka {
5960
@Getter
6061
public static class RocketMQ {
6162

63+
private String topic;
64+
65+
private String namespace;
66+
67+
private String tags;
68+
69+
private String keys;
6270
}
6371
}
6472
}

eden-components/eden-solutions/eden-event-auditor/src/main/java/org/ylzl/eden/event/auditor/integration/kafka/KafkaEventSender.java

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,26 +16,68 @@
1616

1717
package org.ylzl.eden.event.auditor.integration.kafka;
1818

19+
import lombok.RequiredArgsConstructor;
20+
import lombok.extern.slf4j.Slf4j;
21+
import org.springframework.kafka.core.KafkaTemplate;
22+
import org.springframework.kafka.support.SendResult;
23+
import org.springframework.util.concurrent.ListenableFuture;
24+
import org.springframework.util.concurrent.ListenableFutureCallback;
1925
import org.ylzl.eden.event.auditor.EventSender;
2026
import org.ylzl.eden.event.auditor.model.AuditingEvent;
2127

2228
import java.util.List;
29+
import java.util.stream.Collectors;
2330

2431
/**
2532
* 基于 Kafka 发送审计事件
2633
*
2734
* @author <a href="mailto:shiyindaxiaojie@gmail.com">gyl</a>
2835
* @since 2.4.x
2936
*/
37+
@RequiredArgsConstructor
38+
@Slf4j
3039
public class KafkaEventSender implements EventSender {
3140

41+
private static final String KAFKA_SEND_AUDIT_EVENT_SUCCESS = "KafkaTemplate send audit event success, message: {}";
42+
43+
private static final String KAFKA_SEND_AUDIT_EVENT_FAILED = "KafkaTemplate send audit event failed, message: {}";
44+
45+
private static final String KAFKA_SEND_AUDIT_EVENT_ERROR = "KafkaTemplate send audit event error";
46+
47+
private final KafkaTemplate<String, String> kafkaTemplate;
48+
49+
private final String topic;
50+
3251
/**
3352
* 发送审计事件列表
3453
*
3554
* @param events 审计事件列表
3655
*/
3756
@Override
3857
public void send(List<AuditingEvent> events) {
58+
List<String> messages = events.stream()
59+
.map(AuditingEvent::getContent).collect(Collectors.toList());
60+
messages.forEach(
61+
message -> {
62+
try {
63+
ListenableFuture<SendResult<String, String>> future =
64+
kafkaTemplate.send(topic, message);
65+
future.addCallback(new ListenableFutureCallback<SendResult<String, String>>() {
66+
67+
@Override
68+
public void onSuccess(SendResult<String, String> sendResult) {
69+
log.debug(KAFKA_SEND_AUDIT_EVENT_SUCCESS, message);
70+
}
3971

72+
@Override
73+
public void onFailure(Throwable e) {
74+
log.warn(KAFKA_SEND_AUDIT_EVENT_FAILED, message, e);
75+
}
76+
});
77+
} catch (Exception e) {
78+
log.error(KAFKA_SEND_AUDIT_EVENT_ERROR, e);
79+
}
80+
}
81+
);
4082
}
4183
}

eden-components/eden-solutions/eden-event-auditor/src/main/java/org/ylzl/eden/event/auditor/integration/kafka/KafkaEventSenderBuilder.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,9 +18,12 @@
1818
package org.ylzl.eden.event.auditor.integration.kafka;
1919

2020
import lombok.extern.slf4j.Slf4j;
21+
import org.springframework.kafka.core.KafkaTemplate;
2122
import org.ylzl.eden.event.auditor.EventSender;
2223
import org.ylzl.eden.event.auditor.builder.AbstractEventSenderBuilder;
2324
import org.ylzl.eden.event.auditor.builder.EventSenderBuilder;
25+
import org.ylzl.eden.event.auditor.config.EventAuditorConfig;
26+
import org.ylzl.eden.spring.framework.beans.ApplicationContextHelper;
2427

2528
/**
2629
* 基于 Kafka 发送审计事件构建器
@@ -36,8 +39,11 @@ public class KafkaEventSenderBuilder extends AbstractEventSenderBuilder implemen
3639
*
3740
* @return 事件审计实例
3841
*/
42+
@SuppressWarnings("unchecked")
3943
@Override
4044
public EventSender build() {
41-
return null;
45+
KafkaTemplate<String, String> kafkaTemplate = ApplicationContextHelper.getBean(KafkaTemplate.class);
46+
EventAuditorConfig.Sender.Kafka config = this.getEventAuditorConfig().getSender().getKafka();
47+
return new KafkaEventSender(kafkaTemplate, config.getTopic());
4248
}
4349
}

0 commit comments

Comments
 (0)