| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403 |
- import { writeFileSync } from 'node:fs';
- import { fileURLToPath } from 'node:url';
- import path from 'node:path';
- const base = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '../src/main/java/com/fs/company/service/easycall');
- const helper = `package com.fs.company.service.easycall;
- import com.fs.common.core.redis.RedisCacheT;
- import org.springframework.beans.factory.annotation.Autowired;
- import org.springframework.data.redis.core.RedisTemplate;
- import org.springframework.stereotype.Component;
- import java.util.ArrayList;
- import java.util.Collections;
- import java.util.List;
- import java.util.Map;
- import java.util.Set;
- import java.util.concurrent.TimeUnit;
- /**
- * ${'\u5916\u547C\u9650\u6D41\u91CD\u8BD5\u961F\u5217\u4E13\u7528 Redis \u64CD\u4F5C\uFF08ZSET \u5F39\u51FA\u3001Hash \u8F7D\u8377\u3001\u5904\u7406\u9501\u3001\u7F51\u5173\u9650\u901F\u7B49\uFF09\u3002'}
- * <p>
- * ${'\u4E1A\u52A1\u8BED\u4E49\u4E0D\u8FDB'} {@link com.fs.common.core.redis.RedisCache} / {@link RedisCacheT} ${'\u901A\u7528\u5DE5\u5177\u7C7B\u3002'}
- * <p>
- * ${'\u4E0D\u4F7F\u7528 Lua \u811A\u672C\uFF1ARedisson \u4F20\u5165 ARGV \u65F6\u53EF\u80FD\u5E26\u5F15\u53F7/\u5B57\u8282\u6570\u7EC4\uFF0C\u89E6\u53D1'} ERR Lua redis() command arguments must be strings or integers。
- */
- @Component
- public class OutboundLimitRetryRedisHelper {
- @Autowired
- private RedisCacheT<String> redisCacheT;
- /**
- * Pop due members via Java ZSet API (Redisson-safe; same as WxAddWxFrequencyRedisHelper).
- */
- @SuppressWarnings({"unchecked", "rawtypes"})
- public List<String> popDueMembersByScore(String zsetKey, double maxScore, int count) {
- if (count <= 0) {
- return Collections.emptyList();
- }
- RedisTemplate redisTemplate = redisCacheT.redisTemplate;
- Set<Object> members = redisTemplate.opsForZSet().rangeByScore(zsetKey, 0, maxScore, 0, count);
- if (members == null || members.isEmpty()) {
- return Collections.emptyList();
- }
- List<String> dueMembers = new ArrayList<>(members.size());
- for (Object member : members) {
- Long removed = redisTemplate.opsForZSet().remove(zsetKey, member);
- if (removed != null && removed > 0) {
- dueMembers.add(String.valueOf(member));
- }
- }
- return dueMembers;
- }
- public Double zScore(String zsetKey, Object member) {
- return redisCacheT.redisTemplate.opsForZSet().score(zsetKey, member);
- }
- public void zSetAdd(String zsetKey, Object member, double score) {
- redisCacheT.redisTemplate.opsForZSet().add(zsetKey, member, score);
- }
- public void hashPut(String hashKey, String field, String value) {
- redisCacheT.setCacheMapValue(hashKey, field, value);
- }
- public String hashGet(String hashKey, String field) {
- return redisCacheT.getCacheMapValue(hashKey, field);
- }
- public Map<String, String> hashGetAll(String hashKey) {
- Map<String, String> map = redisCacheT.getCacheMap(hashKey);
- return map != null ? map : Collections.emptyMap();
- }
- public void hashDelete(String hashKey, String field) {
- redisCacheT.redisTemplate.opsForHash().delete(hashKey, field);
- }
- public boolean setIfAbsent(String key, String value, long timeout, TimeUnit unit) {
- return Boolean.TRUE.equals(
- redisCacheT.redisTemplate.opsForValue().setIfAbsent(key, value, timeout, unit));
- }
- public void deleteKey(String key) {
- redisCacheT.deleteObject(key);
- }
- public String getString(String key) {
- return redisCacheT.getCacheObject(key);
- }
- public void setString(String key, String value, long timeout, TimeUnit unit) {
- redisCacheT.setCacheObject(key, value, timeout, unit);
- }
- public Long incr(String key, long delta) {
- return redisCacheT.redisTemplate.opsForValue().increment(key, delta);
- }
- public boolean expire(String key, long timeout, TimeUnit unit) {
- return redisCacheT.expire(key, timeout, unit);
- }
- }
- `;
- const support = `package com.fs.company.service.easycall;
- import com.alibaba.fastjson.JSON;
- import com.alibaba.fastjson.JSONObject;
- import com.fs.common.utils.StringUtils;
- import com.fs.company.vo.easycall.EasyCallCommonAddCallListParam;
- import com.fs.company.vo.easycall.EasyCallPhoneItemVO;
- import lombok.extern.slf4j.Slf4j;
- import org.springframework.beans.factory.annotation.Autowired;
- import org.springframework.stereotype.Component;
- import java.util.Date;
- import java.util.List;
- import java.util.Map;
- import java.util.UUID;
- import java.util.concurrent.ThreadLocalRandom;
- import java.util.concurrent.TimeUnit;
- /**
- * ${'\u5916\u547C\u7EBF\u8DEF\u9650\u6D41\u91CD\u8BD5\u961F\u5217\uFF1A\u53BB\u91CD\u5165\u961F\u3001\u6296\u52A8\u5206\u6563\u3001\u539F\u5B50\u51FA\u961F\u3001\u5931\u8D25\u9000\u907F\u91CD\u5165\u961F\u3002'}
- */
- @Slf4j
- @Component
- public class OutboundLimitRetrySupport {
- /** ${'\u4E0E EasyCallServiceImpl \u4FDD\u6301\u4E00\u81F4'} */
- public static final String OUTBOUND_LIMIT_REDIS_PREFIX = "outbound:limit:";
- public static final String RETRY_ZSET_KEY = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:zset";
- public static final String RETRY_PAYLOAD_HASH_KEY = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:payload";
- private static final String RETRY_PROCESSING_PREFIX = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:processing:";
- private static final String RETRY_DISPATCHED_PREFIX = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:dispatched:";
- /** ${'\u5DF2\u6210\u529F append+startTask \u7684 callBackUuid \u6807\u8BB0 TTL\uFF08\u9632 pop \u540E\u5D29\u6E83\u5BFC\u81F4\u91CD\u590D\u5916\u547C\uFF09'} */
- private static final int DISPATCHED_TTL_HOURS = 24;
- private static final String RETRY_RATE_PREFIX = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:rate:";
- /** ${'\u5230\u671F\u540E\u968F\u673A\u5EF6\u8FDF\u4E0A\u9650\uFF08\u6BEB\u79D2\uFF09\uFF0C\u5E73\u6ED1\u7A97\u53E3\u8FB9\u754C\u6D41\u91CF'} */
- private static final int RETRY_JITTER_MAX_MS = 30_000;
- /** ${'\u5B64\u513F payload \u5355\u6B21\u6700\u591A\u6062\u590D\u6761\u6570'} */
- private static final int MAX_ORPHAN_RECOVER_BATCH = 100;
- /** ${'\u5931\u8D25\u9000\u907F\u57FA\u7840\u5EF6\u8FDF\uFF08\u6BEB\u79D2\uFF09'} */
- private static final long RETRY_BACKOFF_MS = 60_000L;
- /** ${'\u5355\u7F51\u5173\u6BCF\u79D2\u6700\u591A\u91CD\u8BD5\u5916\u547C\u6B21\u6570\uFF08\u591A\u5B9E\u4F8B\u5171\u4EAB\uFF09'} */
- private static final int MAX_RETRY_PER_GATEWAY_PER_SEC = 30;
- /** ${'\u5904\u7406\u9501 TTL\uFF0C\u9632\u6B62\u5B9E\u4F8B\u5D29\u6E83\u540E\u957F\u671F\u5360\u9501'} */
- private static final int PROCESSING_LOCK_SECONDS = 600;
- @Autowired
- private OutboundLimitRetryRedisHelper retryRedisHelper;
- /**
- * ${'\u9650\u6D41\u540E\u5165\u961F\uFF08member=callBackUuid\uFF0C\u540C UUID \u4E0D\u91CD\u590D\u5165\u961F\uFF0Cscore \u5E26\u6296\u52A8\uFF09\u3002'}
- */
- public void enqueueRetry(Long companyId, Long gatewayId, EasyCallCommonAddCallListParam param, Date nextAvailableTime) {
- if (param == null || nextAvailableTime == null) {
- return;
- }
- try {
- String callBackUuid = extractCallBackUuid(param);
- String member = StringUtils.isNotBlank(callBackUuid) ? callBackUuid : UUID.randomUUID().toString();
- JSONObject retryData = new JSONObject();
- retryData.put("companyId", companyId);
- retryData.put("gatewayId", gatewayId);
- retryData.put("param", param);
- retryData.put("callBackUuid", member);
- retryData.put("nextAvailableTime", nextAvailableTime.getTime());
- retryData.put("createTime", System.currentTimeMillis());
- String existingPayload = retryRedisHelper.hashGet(RETRY_PAYLOAD_HASH_KEY, member);
- if (existingPayload != null) {
- JSONObject oldData = JSON.parseObject(existingPayload);
- if (oldData != null && oldData.containsKey("failCount")) {
- retryData.put("failCount", oldData.getIntValue("failCount"));
- }
- }
- double score = computeScoreWithJitter(nextAvailableTime.getTime());
- Double existingScore = retryRedisHelper.zScore(RETRY_ZSET_KEY, member);
- if (existingScore != null && existingScore <= score) {
- retryRedisHelper.hashPut(RETRY_PAYLOAD_HASH_KEY, member, retryData.toJSONString());
- log.info("enqueueRetry: ${'\u5DF2\u5B58\u5728\u66F4\u65E9\u6216\u76F8\u540C\u8C03\u5EA6\uFF0C\u4EC5\u5237\u65B0 payload'} - member={}", member);
- return;
- }
- retryRedisHelper.zSetAdd(RETRY_ZSET_KEY, member, score);
- retryRedisHelper.hashPut(RETRY_PAYLOAD_HASH_KEY, member, retryData.toJSONString());
- log.info("enqueueRetry: member={}, score={}, companyId={}, gatewayId={}", member, (long) score, companyId, gatewayId);
- } catch (Exception e) {
- log.error("enqueueRetry ${'\u5F02\u5E38'} companyId={}, gatewayId={}", companyId, gatewayId, e);
- }
- }
- /**
- * ${'\u5F39\u51FA\u5DF2\u5230\u671F\u7684\u91CD\u8BD5 member\uFF08Java ZSet rangeByScore + remove\uFF0C\u4E0E\u52A0\u5FAE\u9891\u7387\u91CD\u8BD5\u540C\u6A21\u5F0F\uFF09\u3002'}
- */
- public List<String> popDueMembers(int batchSize) {
- long now = System.currentTimeMillis();
- return retryRedisHelper.popDueMembersByScore(RETRY_ZSET_KEY, now, batchSize);
- }
- /**
- * ${'\u89E3\u6790\u4EFB\u52A1 payload\uFF1B\u517C\u5BB9\u65E7\u7248 member \u4E3A\u6574\u6BB5 JSON \u7684\u5199\u6CD5\u3002'}
- */
- public JSONObject resolvePayload(String member) {
- if (StringUtils.isBlank(member)) {
- return null;
- }
- if (member.trim().startsWith("{")) {
- return JSON.parseObject(member);
- }
- String raw = retryRedisHelper.hashGet(RETRY_PAYLOAD_HASH_KEY, member);
- if (raw == null) {
- return null;
- }
- return JSON.parseObject(raw);
- }
- /**
- * ${'\u5C1D\u8BD5\u83B7\u53D6\u5904\u7406\u9501\uFF0C\u907F\u514D\u6781\u7AEF\u60C5\u51B5\u4E0B\u91CD\u590D\u5916\u547C\u3002'}
- */
- public boolean tryAcquireProcessingLock(String callBackUuid) {
- if (StringUtils.isBlank(callBackUuid)) {
- return true;
- }
- String lockKey = RETRY_PROCESSING_PREFIX + callBackUuid;
- return retryRedisHelper.setIfAbsent(lockKey, "1", PROCESSING_LOCK_SECONDS, TimeUnit.SECONDS);
- }
- public void releaseProcessingLock(String callBackUuid) {
- if (StringUtils.isBlank(callBackUuid)) {
- return;
- }
- retryRedisHelper.deleteKey(RETRY_PROCESSING_PREFIX + callBackUuid);
- }
- /**
- * ${'\u662F\u5426\u5DF2\u6210\u529F\u5B8C\u6210 append+startTask\uFF08\u9632 pop \u540E\u8FDB\u7A0B\u5D29\u6E83\u518D\u6B21\u91CD\u8BD5\u5BFC\u81F4\u91CD\u590D\u5916\u547C\uFF09\u3002'}
- */
- public boolean isDispatched(String callBackUuid) {
- if (StringUtils.isBlank(callBackUuid)) {
- return false;
- }
- return retryRedisHelper.getString(RETRY_DISPATCHED_PREFIX + callBackUuid) != null;
- }
- /**
- * ${'\u6807\u8BB0 callBackUuid \u5DF2\u5B8C\u6210\u6D3E\u53D1\uFF08append \u6210\u529F\u4E14 startTask \u5DF2\u89E6\u53D1\uFF09\u3002'}
- */
- public void markDispatched(String callBackUuid) {
- if (StringUtils.isBlank(callBackUuid)) {
- return;
- }
- retryRedisHelper.setString(RETRY_DISPATCHED_PREFIX + callBackUuid, "1", DISPATCHED_TTL_HOURS, TimeUnit.HOURS);
- }
- /**
- * ${'\u6062\u590D ZSET \u5DF2\u5F39\u51FA\u4F46 payload \u4ECD\u6EDE\u7559 Hash \u7684\u5B64\u513F\u4EFB\u52A1\uFF08\u8FDB\u7A0B\u5D29\u6E83\u573A\u666F\uFF0C\u4E0D\u91CD\u8BD5\u4E22\u5355\uFF09\u3002'}
- *
- * @return ${'\u672C\u6B21\u6062\u590D\u5165\u961F\u6761\u6570'}
- */
- public int recoverOrphanedPayloads() {
- Map<String, String> all = retryRedisHelper.hashGetAll(RETRY_PAYLOAD_HASH_KEY);
- if (all == null || all.isEmpty()) {
- return 0;
- }
- int recovered = 0;
- for (Map.Entry<String, String> entry : all.entrySet()) {
- if (recovered >= MAX_ORPHAN_RECOVER_BATCH) {
- break;
- }
- String member = entry.getKey();
- if (retryRedisHelper.zScore(RETRY_ZSET_KEY, member) != null) {
- continue;
- }
- JSONObject data = JSON.parseObject(entry.getValue());
- if (data == null) {
- ackSuccess(member);
- continue;
- }
- String callBackUuid = data.getString("callBackUuid");
- EasyCallCommonAddCallListParam param = data.getObject("param", EasyCallCommonAddCallListParam.class);
- if (StringUtils.isBlank(callBackUuid) && param != null) {
- callBackUuid = extractCallBackUuid(param);
- }
- if (isDispatched(callBackUuid)) {
- ackSuccess(member);
- continue;
- }
- Long companyId = data.getLong("companyId");
- Long gatewayId = data.getLong("gatewayId");
- if (companyId == null || param == null) {
- ackSuccess(member);
- continue;
- }
- Date nextTime = new Date(System.currentTimeMillis() + RETRY_BACKOFF_MS);
- enqueueRetry(companyId, gatewayId, param, nextTime);
- recovered++;
- log.warn("recoverOrphanedPayloads: ${'\u5B64\u513F\u4EFB\u52A1\u91CD\u65B0\u5165\u961F'} member={}, callBackUuid={}", member, callBackUuid);
- }
- return recovered;
- }
- /**
- * ${'\u6309\u7F51\u5173\u4EE4\u724C\u6876\u9650\u901F\uFF0C\u8D85\u9650\u5219\u5EF6\u8FDF\u91CD\u5165\u961F\u3002'}
- *
- * @return true ${'\u8868\u793A\u5F53\u524D\u79D2\u989D\u5EA6\u5DF2\u7528\u5C3D\uFF0C\u8C03\u7528\u65B9\u5E94\u8DF3\u8FC7\u672C\u6B21\u5916\u547C'}
- */
- public boolean isGatewayRateLimited(Long companyId, Long gatewayId) {
- if (companyId == null || gatewayId == null) {
- return false;
- }
- String bucketKey = RETRY_RATE_PREFIX + companyId + ":" + gatewayId;
- Long count = retryRedisHelper.incr(bucketKey, 1L);
- if (count == null) {
- count = 1L;
- }
- if (count == 1L) {
- retryRedisHelper.expire(bucketKey, 1, TimeUnit.SECONDS);
- }
- return count > MAX_RETRY_PER_GATEWAY_PER_SEC;
- }
- /**
- * ${'\u6210\u529F\u540E\u6E05\u7406 payload\uFF1Bmember \u5DF2\u5728 pop \u65F6\u4ECE ZSET \u79FB\u9664\u3002'}
- */
- public void ackSuccess(String member) {
- if (StringUtils.isBlank(member) || member.trim().startsWith("{")) {
- return;
- }
- retryRedisHelper.hashDelete(RETRY_PAYLOAD_HASH_KEY, member);
- }
- /**
- * ${'\u5931\u8D25\u6216\u5F02\u5E38\uFF1A\u9000\u907F\u540E\u91CD\u65B0\u5165\u961F\uFF0C\u4E0D\u4E22\u5355\u3002'}
- */
- public void requeueWithBackoff(JSONObject retryData, String member) {
- if (retryData == null) {
- return;
- }
- Long companyId = retryData.getLong("companyId");
- Long gatewayId = retryData.getLong("gatewayId");
- EasyCallCommonAddCallListParam param = retryData.getObject("param", EasyCallCommonAddCallListParam.class);
- if (param == null) {
- log.warn("requeueWithBackoff: param ${'\u4E3A\u7A7A\uFF0C\u65E0\u6CD5\u91CD\u5165\u961F'} member={}", member);
- return;
- }
- int failCount = retryData.getIntValue("failCount") + 1;
- retryData.put("failCount", failCount);
- if (StringUtils.isNotBlank(member) && !member.trim().startsWith("{")) {
- retryRedisHelper.hashPut(RETRY_PAYLOAD_HASH_KEY, member, retryData.toJSONString());
- }
- long delay = RETRY_BACKOFF_MS * Math.min(failCount, 5);
- Date nextTime = new Date(System.currentTimeMillis() + delay);
- enqueueRetry(companyId, gatewayId, param, nextTime);
- log.warn("requeueWithBackoff: member={}, failCount={}, nextDelayMs={}", member, failCount, delay);
- }
- public static String extractCallBackUuid(EasyCallCommonAddCallListParam param) {
- if (param == null || param.getPhoneList() == null || param.getPhoneList().isEmpty()) {
- return null;
- }
- EasyCallPhoneItemVO item = param.getPhoneList().get(0);
- if (item == null || item.getBizJson() == null) {
- return null;
- }
- JSONObject bizJson;
- if (item.getBizJson() instanceof JSONObject) {
- bizJson = (JSONObject) item.getBizJson();
- } else {
- bizJson = JSON.parseObject(String.valueOf(item.getBizJson()));
- }
- return bizJson != null ? bizJson.getString("callBackUuid") : null;
- }
- private static double computeScoreWithJitter(long baseTimeMs) {
- int jitter = ThreadLocalRandom.current().nextInt(RETRY_JITTER_MAX_MS + 1);
- return baseTimeMs + jitter;
- }
- }
- `;
- writeFileSync(path.join(base, 'OutboundLimitRetryRedisHelper.java'), helper, 'utf8');
- writeFileSync(path.join(base, 'OutboundLimitRetrySupport.java'), support, 'utf8');
- console.log('UTF-8 files written');
|