yumeng 1 неделя назад
Родитель
Сommit
9c93feb63d

+ 7 - 12
src/main/java/com/adx/tencent/conversionsync/ConversionSyncService.java

@@ -46,7 +46,6 @@ public class ConversionSyncService {
     private static final Logger log = LoggerFactory.getLogger(ConversionSyncService.class);
     private static final int REALTIME_CONCURRENCY = 10;
     private static final int BACKFILL_CONCURRENCY = 16;
-    private static final int DEDUCTION_WINDOW_SIZE = 100;
 
     private final ConversionClient baiduClient;
     private final RedisHotStore hotStore;
@@ -592,17 +591,18 @@ public class ConversionSyncService {
                     task.callbackKey, task.payment.getQk(), safe(task.tencentBid.getAccountId()));
             return false;
         }
-        long counter = hotStore.incrementDeductionCounter(
+        RedisHotStore.DeductionDecision decision = hotStore.evaluateDeduction(
                 task.tencentBid.getAccountId(),
                 task.payment.getAct(),
                 resolveTrackingVersion(task.tencentBid),
-                task.payment.getDate()
+                task.payment.getDate(),
+                deductionRate
         );
-        int slot = deductionSlot(counter);
-        boolean deduct = slot < deductionRate;
-        log.info("[ConversionSync] 腾讯转化扣量判定 | qk={} | act={} | accountId={} | version={} | rate={} | counter={} | slot={} | callbackKey={} | deduct={}",
+        boolean deduct = decision.isShouldDeduct();
+        log.info("[ConversionSync] 腾讯转化扣量判定 | qk={} | act={} | accountId={} | version={} | rate={} | totalSeen={} | totalDeducted={} | callbackKey={} | deduct={}",
                 task.payment.getQk(), task.payment.getAct(), safe(task.tencentBid.getAccountId()),
-                resolveTrackingVersion(task.tencentBid), deductionRate, counter, slot, task.callbackKey, deduct);
+                resolveTrackingVersion(task.tencentBid), deductionRate,
+                decision.getTotalSeen(), decision.getTotalDeducted(), task.callbackKey, deduct);
         return deduct;
     }
 
@@ -647,11 +647,6 @@ public class ConversionSyncService {
         return record;
     }
 
-    private int deductionSlot(long counter) {
-        long normalized = Math.max(counter - 1L, 0L);
-        return (int) (normalized % DEDUCTION_WINDOW_SIZE);
-    }
-
     private static String resolveTrackingVersion(BidRecord bid) {
         if (bid == null || bid.getMediaParams() == null) {
             return "v1";

+ 73 - 10
src/main/java/com/adx/tencent/storage/RedisHotStore.java

@@ -4,6 +4,7 @@ import com.adx.tencent.storage.model.BidRecord;
 import com.adx.tencent.storage.model.QueuedEvent;
 import com.adx.tencent.storage.model.TrackingRecord;
 import com.fasterxml.jackson.databind.ObjectMapper;
+import org.springframework.data.redis.core.script.DefaultRedisScript;
 import org.springframework.data.redis.connection.stream.*;
 import org.springframework.data.redis.core.StringRedisTemplate;
 
@@ -35,6 +36,7 @@ public class RedisHotStore {
     private final String prefix;
     private final String stream;
     private final Duration bidTtl;
+    private final DefaultRedisScript<List> deductionDecisionScript;
     private final Set<String> initializedGroups = ConcurrentHashMap.newKeySet();
 
     public RedisHotStore(StringRedisTemplate redis, ObjectMapper objectMapper,
@@ -45,6 +47,7 @@ public class RedisHotStore {
         this.prefix = (prefix == null || prefix.isBlank()) ? DEFAULT_PREFIX : prefix;
         this.stream = (stream == null || stream.isBlank()) ? DEFAULT_STREAM : stream;
         this.bidTtl = (bidTtl == null || bidTtl.isZero()) ? Duration.ofHours(24) : bidTtl;
+        this.deductionDecisionScript = buildDeductionDecisionScript();
     }
 
     // ─── EventRecorder ───────────────────────────────────────────────────────
@@ -221,18 +224,24 @@ public class RedisHotStore {
     }
 
     /**
-     * 按账户/事件/版本/日期维度维护一个递增计数器,供窗口扣量使用。
+     * 原子计算一条新转化是否应被扣量。
+     * 规则:任意时点都尽量逼近目标比例,例如 50 条时扣 25,100 条时扣 50。
      */
-    public long incrementDeductionCounter(String accountId, int act, String trackingVersion, String date) {
+    public DeductionDecision evaluateDeduction(String accountId, int act, String trackingVersion, String date, int deductionRate) {
         String key = deductionCounterKey(accountId, act, trackingVersion, date);
-        Long count = redis.opsForValue().increment(key);
-        if (count == null) {
-            throw new RuntimeException("increment deduction counter failed");
-        }
-        if (count == 1L) {
-            redis.expire(key, DEFAULT_DEDUCTION_COUNTER_TTL);
+        List result = redis.execute(
+                deductionDecisionScript,
+                Collections.singletonList(key),
+                String.valueOf(Math.max(deductionRate, 0)),
+                String.valueOf(DEFAULT_DEDUCTION_COUNTER_TTL.getSeconds())
+        );
+        if (result == null || result.size() < 3) {
+            throw new RuntimeException("evaluate deduction failed");
         }
-        return count;
+        long total = toLong(result.get(0));
+        long deducted = toLong(result.get(1));
+        boolean shouldDeduct = toLong(result.get(2)) == 1L;
+        return new DeductionDecision(total, deducted, shouldDeduct);
     }
 
     // ─── EventQueue ──────────────────────────────────────────────────────────
@@ -375,7 +384,7 @@ public class RedisHotStore {
     }
 
     private String deductionCounterKey(String accountId, int act, String trackingVersion, String date) {
-        return String.format("%sdeduct:counter:tencent:acct:%s:act:%d:ver:%s:date:%s",
+        return String.format("%sdeduct:state:tencent:acct:%s:act:%d:ver:%s:date:%s",
                 prefix,
                 sanitizeCounterPart(accountId, "unknown"),
                 act,
@@ -390,6 +399,60 @@ public class RedisHotStore {
         return value.replace(':', '_').trim();
     }
 
+    private static long toLong(Object value) {
+        if (value instanceof Number number) {
+            return number.longValue();
+        }
+        return Long.parseLong(String.valueOf(value));
+    }
+
+    private DefaultRedisScript<List> buildDeductionDecisionScript() {
+        DefaultRedisScript<List> script = new DefaultRedisScript<>();
+        script.setResultType(List.class);
+        script.setScriptText("""
+                local key = KEYS[1]
+                local rate = tonumber(ARGV[1]) or 0
+                local ttl = tonumber(ARGV[2]) or 0
+                local total = redis.call('HINCRBY', key, 'total', 1)
+                local deducted = tonumber(redis.call('HGET', key, 'deducted') or '0')
+                local target = math.floor((total * rate) / 100)
+                local should_deduct = 0
+                if deducted < target then
+                    deducted = redis.call('HINCRBY', key, 'deducted', 1)
+                    should_deduct = 1
+                end
+                if ttl > 0 then
+                    redis.call('EXPIRE', key, ttl)
+                end
+                return { total, deducted, should_deduct }
+                """);
+        return script;
+    }
+
+    public static class DeductionDecision {
+        private final long totalSeen;
+        private final long totalDeducted;
+        private final boolean shouldDeduct;
+
+        public DeductionDecision(long totalSeen, long totalDeducted, boolean shouldDeduct) {
+            this.totalSeen = totalSeen;
+            this.totalDeducted = totalDeducted;
+            this.shouldDeduct = shouldDeduct;
+        }
+
+        public long getTotalSeen() {
+            return totalSeen;
+        }
+
+        public long getTotalDeducted() {
+            return totalDeducted;
+        }
+
+        public boolean isShouldDeduct() {
+            return shouldDeduct;
+        }
+    }
+
     private BidRecord parseBidRecord(String payload) {
         try {
             return objectMapper.readValue(payload, BidRecord.class);