From 34981c30a78e8bbd7791131059a9210f9928b62c Mon Sep 17 00:00:00 2001
From: 刘光辉 <347230014@qq.com>
Date: 星期四, 17 九月 2026 09:24:09 +0800
Subject: [PATCH] Merge remote-tracking branch 'origin/master' into master
---
jnpf-biz-common/jnpf-audit-sdk/src/main/java/jnpf/audit/sdk/AuditEventPublisher.java | 131 +++++++++++++++++++++++++++++++++++++++++++
1 files changed, 131 insertions(+), 0 deletions(-)
diff --git a/jnpf-biz-common/jnpf-audit-sdk/src/main/java/jnpf/audit/sdk/AuditEventPublisher.java b/jnpf-biz-common/jnpf-audit-sdk/src/main/java/jnpf/audit/sdk/AuditEventPublisher.java
new file mode 100644
index 0000000..d63bfe4
--- /dev/null
+++ b/jnpf-biz-common/jnpf-audit-sdk/src/main/java/jnpf/audit/sdk/AuditEventPublisher.java
@@ -0,0 +1,131 @@
+package jnpf.audit.sdk;
+
+import jnpf.audit.AuditDeliveryException;
+import jnpf.audit.AuditEventApi;
+import jnpf.audit.model.AuditEventDTO;
+import jnpf.base.ActionResult;
+import jnpf.base.ActionResultCode;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.DisposableBean;
+import org.springframework.beans.factory.ObjectProvider;
+
+import java.util.List;
+import java.util.Set;
+import java.util.concurrent.ArrayBlockingQueue;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * 涓ょ骇闄嶇骇寮傛鍙戦�佸櫒锛欶eign锛圓uditEventApi 鍚屾锛夆啋 鏈湴 spool 钀界洏銆�
+ *
+ * <p><b>2026-08-11 鍙樻洿</b>锛氬師绗竴绾� MQ锛坽@code streamBridge.send("AUDIT_EVENT", 鈥�)}锛夊凡绉婚櫎銆�
+ * 浜嬩欢鎬荤嚎浜� 2026-08-10 鍒� Redis 鍚庢爤鍐呬笉鍐嶆湁 RocketMQ binder锛岃绾�**姣忔蹇呯劧澶辫触**鍐嶇敱 Feign 鎺ョ鈥斺��
+ * 瀹冨凡涓嶆槸闄嶇骇绾ф锛屽彧鏄瘡鏉′簨浠跺涓�娆″紓甯告瀯閫犮�傜Щ闄ゅ悗琛屼负绛変环锛堟鍓嶆湰灏辨瘡鏉¢兘璧� Feign锛夈��
+ * 鑻ュ皢鏉ヨ鎭㈠寮傛鎶曢�掞紝姝hВ鏄彁渚涗竴涓啓 Redis 鐨� {@link AuditEventApi} 瀹炵幇 bean 鏇夸唬 Feign锛�
+ * 鏈被涓�琛屼笉鐢ㄦ敼锛堣 tasks/todo.md 闃舵 4锛屽綋鍓嶅洜 Redis 鏃犳寔涔呭寲 + noeviction 鏆備笉鍋氾級銆�
+ *
+ * <p>鏈被鐨勫苟鍙�/鏃跺簭娉ㄩ噴瀵瑰簲闄嶇骇閾捐矾涓婄粡杩囬獙璇佺殑骞跺彂瀹夊叏璁捐鈥斺�旈噸鏋勬椂涓嶈闅忔墜鍒犳敼锛�
+ * 姣忎竴澶勯兘瀵瑰簲涓�涓湡瀹炲瓨鍦ㄨ繃鐨勭珵鎬佸満鏅��
+ * deliverOnce 鍙仛 Feign 涓�绾у苟缁欏嚭鏄庣‘鎴愯触淇″彿锛堜笉鍐� spool锛夛紝渚涙櫘閫氳矾寰勪笌 replayer 鍚勮嚜鍐冲畾澶辫触澶勭悊銆�
+ */
+@Slf4j
+public class AuditEventPublisher implements DisposableBean {
+ private final ObjectProvider<AuditEventApi> feignApi;
+ private final AuditSpoolWriter spool;
+ private final ThreadPoolExecutor pool = new ThreadPoolExecutor(1, 2, 60, TimeUnit.SECONDS,
+ new ArrayBlockingQueue<>(1000),
+ r -> { Thread t = new Thread(r, "audit-sender"); t.setDaemon(true); return t; },
+ new ThreadPoolExecutor.AbortPolicy()); // 闃熷垪婊¤蛋 catch 鈫� spool锛屼笉闃诲涓氬姟
+
+ /** pending 闆嗗悎锛氭彁浜ゅ墠鍦ㄨ皟鐢ㄧ嚎绋嬬櫥璁般�乺un() finally 绉婚櫎鈥斺�旂粺涓�瑕嗙洊"鎺掗槦涓�/鎵ц涓�/宸插嚭闃熸湭寮�璺�"涓夋�侊紝
+ 鏃犵櫥璁扮獥鍙o紙缁堥獙涓夎疆鈶o細鍘� run() 鍐� add 鐨勫啓娉曞湪鍑洪槦涓庣櫥璁颁箣闂存湁绐楀彛锛宻hutdownNow 鎾炰笂浼氬弻婕忥級 */
+ private final Set<AuditSendTask> pending = ConcurrentHashMap.newKeySet();
+
+ public AuditEventPublisher(ObjectProvider<AuditEventApi> feignApi,
+ AuditSpoolWriter spool) {
+ this.feignApi = feignApi;
+ this.spool = spool;
+ }
+
+ /** 浠诲姟鑷惡浜嬩欢瀵硅薄锛氬紓甯�/鍋滄満 drain 鏃惰兘鎭㈠鍑哄叿浣撲簨浠惰惤 spool锛堣瘎瀹�#2锛岃8 lambda 鎭㈠涓嶅嚭浜嬩欢锛� */
+ class AuditSendTask implements Runnable {
+ final AuditEventDTO event;
+ AuditSendTask(AuditEventDTO event) { this.event = event; }
+ @Override public void run() {
+ try {
+ deliverOnce(event);
+ } catch (Throwable t) { // 鏈�澶栧眰鍏滃簳锛氫换浣曞紓甯革紙鍚� JSON 搴忓垪鍖栵級閮借惤 spool锛岀粷涓嶉潤榛樹涪锛堣瘎瀹�#2锛�
+ spool.write(event);
+ log.error("AUDIT-LOST clientEventId={} type={} target={} 宸茶惤 spool",
+ event.getClientEventId(), event.getEventType(), event.getTargetTable(), t);
+ } finally {
+ pending.remove(this);
+ }
+ }
+ }
+
+ public void publishAsync(List<AuditEventDTO> events) {
+ for (AuditEventDTO e : events) {
+ AuditSendTask task = new AuditSendTask(e);
+ pending.add(task); // 鍏堢櫥璁板啀鍏ユ睜锛堢粓楠屼笁杞懀锛�
+ try {
+ pool.execute(task);
+ } catch (RejectedExecutionException full) {
+ // 鍏堝啓 spool 鍐嶇Щ闄ょ櫥璁帮紙缁堥獙鍥涜疆鈶★細鍙嶅簭浼氬湪 remove 涓� write 涔嬮棿鐣欎笅鍙屼晶涓嶅彲瑙佺獥鍙o級
+ try {
+ spool.write(e); // 闃熷垪婊$洿鎺ヨ惤 spool锛坰pec 搂5.5锛�
+ } finally {
+ pending.remove(task);
+ }
+ }
+ }
+ }
+
+ /**
+ * 鍗曟鎶曢�掞細Feign 鍚屾锛屾垚鍔熻繑鍥炪�佸け璐ユ姏寮傚父鈥斺�旀湰鏂规硶涓嶅啓 spool銆�
+ * 缁堥獙#2/#3锛氭櫘閫氳矾寰勪笌 replayer 閮介渶瑕佹槑纭殑鎴愯触淇″彿鍚勮嚜鍐冲畾澶辫触澶勭悊锛�
+ * 鍘� send() 澶辫触鏃惰嚜宸卞啓 spool 姝e父杩斿洖锛宺eplayer 鏃犳硶鍒ゆ柇璇ヤ笉璇ヤ繚鐣欏緟閲嶈瘯琛屻��
+ */
+ void deliverOnce(AuditEventDTO e) {
+ ActionResult<Integer> result = feignApi.getObject().record(e);
+ if (result == null || !ActionResultCode.SUCCESS.getCode().equals(result.getCode())) {
+ // HTTP 灞傚け璐ョ敱 fallbackFactory 鎶涘嚭锛涗笟鍔″眰澶辫触(HTTP 200 + fail body)闈犳湰鏍¢獙鎹曡幏锛圕odex 鎵规涓� Critical#1锛�
+ throw new AuditDeliveryException("Feign 鎶曢�掕繑鍥為潪鎴愬姛: " + (result == null ? "null" : result.getCode()), null);
+ }
+ }
+
+ /**
+ * Codex 鎵规浜� #5锛氫互 {@link DisposableBean#destroy()} 鏇夸唬 jakarta {@code @PreDestroy}锛�
+ * 璺� boot2(JDK8)/boot3 涓珛鈥斺�攞@code @PreDestroy} 鍦� boot2 璧� {@code javax.annotation}銆�
+ * boot3 璧� {@code jakarta.annotation}锛孲DK 浜岃繘鍒舵棤娉曞悓鏃跺吋瀹癸紱DisposableBean 鏄� Spring 鍘熺敓鎺ュ彛锛�
+ * 涓ょ増涓�鑷淬�係pring 瀵� {@code @Bean} 浼氳嚜鍔ㄨ瘑鍒� DisposableBean锛寋@code AuditSdkAutoConfiguration} 鏃犻渶鏀广��
+ */
+ @Override
+ public void destroy() {
+ drain();
+ }
+
+ public void drain() { // 鍋滄満 drain锛坰pec 搂5.5锛夛紱涓嶅悜澶栨姏寮傚父锛堢粓楠屼笁杞懀锛氬紓甯搁�冨嚭浼氱粫杩囧叏閮ㄥ厹搴曪級
+ pool.shutdown();
+ try {
+ if (!pool.awaitTermination(10, TimeUnit.SECONDS)) {
+ pool.shutdownNow();
+ spoolPending();
+ }
+ } catch (InterruptedException ie) {
+ pool.shutdownNow();
+ spoolPending();
+ Thread.currentThread().interrupt(); // 鍏滃簳鍋氬畬鍐嶆仮澶嶄腑鏂爣璁�
+ }
+ }
+
+ private void spoolPending() {
+ // pending 缁熶竴瑕嗙洊鎺掗槦涓�/鎵ц涓�/鍑洪槦鏈紑璺戜笁鎬侊紱涓庝换鍔¤嚜韬� catch 钀界洏鍙兘閲嶅鈥斺��
+ // spool 閲嶆斁涓庤惤搴撳潎鎸� client_event_id 骞傜瓑锛岄噸澶嶆棤瀹�
+ for (AuditSendTask task : pending) {
+ spool.write(task.event);
+ }
+ }
+}
--
Gitblit v1.8.0