package com.adx.tencent.rta; import com.adx.tencent.rta.model.RtaTagOrderRecord; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.data.redis.connection.RedisConnection; import org.springframework.data.redis.connection.RedisStringCommands; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.data.redis.core.types.Expiration; import java.nio.charset.StandardCharsets; import java.time.Duration; import java.util.LinkedHashMap; import java.util.LinkedHashSet; import java.util.List; import java.util.Map; import java.util.Set; import java.util.stream.Collectors; /** * RTA order_id 同步服务。 * 定时从 rta_tag_orders 表读取全部记录,同步到 Redis: * key = {prefix}{tag_id} * value = order_id 列表(逗号分隔) */ public class RtaOrderSyncService { private static final Logger log = LoggerFactory.getLogger(RtaOrderSyncService.class); private final RtaOrderStore store; private final StringRedisTemplate redis; private final String keyPrefix; private final Duration ttl; public RtaOrderSyncService(RtaOrderStore store, StringRedisTemplate redis, String keyPrefix, Duration ttl) { this.store = store; this.redis = redis; this.keyPrefix = keyPrefix == null ? "adx:rta:order:" : keyPrefix; this.ttl = ttl; } public int syncOnce() { List rows = store.listAll(); if (rows == null || rows.isEmpty()) { log.info("[RtaOrderSync] 无数据可同步"); return 0; } Map> grouped = new LinkedHashMap<>(); for (RtaTagOrderRecord row : rows) { if (row.getTagId() == null || row.getTagId().isBlank()) { continue; } if (row.getOrderId() == null || row.getOrderId().isBlank()) { continue; } grouped.computeIfAbsent(row.getTagId().trim(), ignored -> new LinkedHashSet<>()) .add(row.getOrderId().trim()); } if (grouped.isEmpty()) { log.info("[RtaOrderSync] 无有效数据可同步"); return 0; } long ttlMs = (ttl != null && !ttl.isZero() && !ttl.isNegative()) ? ttl.toMillis() : 0; redis.executePipelined((RedisConnection conn) -> { for (Map.Entry> entry : grouped.entrySet()) { byte[] key = redisKey(entry.getKey()).getBytes(StandardCharsets.UTF_8); String value = entry.getValue().stream().collect(Collectors.joining(",")); byte[] val = value.getBytes(StandardCharsets.UTF_8); if (ttlMs > 0) { conn.stringCommands().set(key, val, Expiration.milliseconds(ttlMs), RedisStringCommands.SetOption.UPSERT); } else { conn.stringCommands().set(key, val); } } return null; }); log.info("[RtaOrderSync] 同步完成,共 {} 条 RTA order_id 写入 Redis, 聚合后 {} 个广告位", rows.size(), grouped.size()); return rows.size(); } private String redisKey(String tagId) { return keyPrefix + tagId; } }