yumeng hace 6 días
padre
commit
0a5484a39b

+ 4 - 2
src/main/java/com/adx/tencent/AppConfiguration.java

@@ -181,7 +181,7 @@ public class AppConfiguration {
     @Bean
     public OppoRetryService oppoRetryService(@Nullable OppoColdStore coldStore, OppoClient client,
                                              ObjectMapper mapper) {
-        return coldStore == null ? null : new OppoRetryService(coldStore, client, mapper);
+        return coldStore == null ? null : new OppoRetryService(coldStore, client, mapper, props.getOppoCallbackRetryLimit());
     }
 
     @Bean
@@ -221,7 +221,9 @@ public class AppConfiguration {
     @Bean
     public OppoColdWorker oppoColdWorker(@Nullable OppoHotStore hotStore, @Nullable OppoColdStore coldStore, ObjectMapper mapper) {
         if (hotStore == null || coldStore == null) return null;
-        return new OppoColdWorker(hotStore, coldStore, "oppo-cold-writers", props.getWorkerConsumer() + "-oppo", props.getWorkerBatch(), props.getRedisStreamMaxLen(), mapper);
+        return new OppoColdWorker(hotStore, coldStore,
+                "oppo-cold-writers", props.getWorkerConsumer() + "-oppo", props.getWorkerBatch(),
+                null, props.getRedisStreamMaxLen(), mapper);
     }
 
     @Bean

+ 2 - 1
src/main/java/com/adx/tencent/oppo/service/OppoBackgroundTasks.java

@@ -78,7 +78,8 @@ public class OppoBackgroundTasks implements ApplicationRunner {
         long intervalMs = props.getOppoCallbackRetryInterval().toMillis();
         while (!stopSignal.isStopped()) {
             try {
-                retryService.retry(props.getOppoCallbackRetryLimit());
+                OppoRetryService.RetryResult r = retryService.retryCallbacks(props.getOppoCallbackRetryLimit());
+                log.info("[OppoCallbackRetry] done: fetched={} sent={} failed={}", r.fetched, r.sent, r.failed);
             } catch (Exception e) {
                 log.error("oppo callback retry: {}", e.getMessage(), e);
             }

La diferencia del archivo ha sido suprimido porque es demasiado grande
+ 57 - 4
src/main/java/com/adx/tencent/oppo/service/OppoColdWorker.java


La diferencia del archivo ha sido suprimido porque es demasiado grande
+ 84 - 5
src/main/java/com/adx/tencent/oppo/service/OppoRetryService.java


La diferencia del archivo ha sido suprimido porque es demasiado grande
+ 72 - 7
src/main/java/com/adx/tencent/oppo/service/OppoTagEventSyncService.java


La diferencia del archivo ha sido suprimido porque es demasiado grande
+ 378 - 34
src/main/java/com/adx/tencent/oppo/store/OppoHotStore.java


+ 35 - 8
src/test/java/com/adx/tencent/oppo/store/OppoHotStoreTest.java

@@ -6,10 +6,15 @@ import com.fasterxml.jackson.databind.ObjectMapper;
 import com.fasterxml.jackson.databind.json.JsonMapper;
 import org.junit.jupiter.api.Test;
 import org.mockito.ArgumentCaptor;
+import org.springframework.data.redis.connection.RedisConnection;
+import org.springframework.data.redis.connection.RedisStringCommands;
+import org.springframework.data.redis.connection.RedisStreamCommands;
 import org.springframework.data.redis.connection.stream.MapRecord;
+import org.springframework.data.redis.core.RedisCallback;
 import org.springframework.data.redis.core.StreamOperations;
 import org.springframework.data.redis.core.StringRedisTemplate;
 import org.springframework.data.redis.core.ValueOperations;
+import org.springframework.data.redis.core.types.Expiration;
 
 import java.time.Duration;
 import java.util.List;
@@ -26,11 +31,21 @@ class OppoHotStoreTest {
     private final ValueOperations<String, String> values = mock(ValueOperations.class);
     @SuppressWarnings("unchecked")
     private final StreamOperations<String, Object, Object> stream = mock(StreamOperations.class);
+    private final RedisConnection connection = mock(RedisConnection.class);
+    private final RedisStringCommands stringCommands = mock(RedisStringCommands.class);
+    private final RedisStreamCommands streamCommands = mock(RedisStreamCommands.class);
     private final OppoColdStore cold = mock(OppoColdStore.class);
 
     private OppoHotStore store(Duration ttl) {
         when(redis.opsForValue()).thenReturn(values);
         when(redis.opsForStream()).thenReturn(stream);
+        when(connection.stringCommands()).thenReturn(stringCommands);
+        when(connection.streamCommands()).thenReturn(streamCommands);
+        when(redis.executePipelined(any(RedisCallback.class))).thenAnswer(invocation -> {
+            RedisCallback<?> callback = invocation.getArgument(0);
+            callback.doInRedis(connection);
+            return List.of();
+        });
         return new OppoHotStore(redis, mapper, cold, "adx:oppo:events", ttl, 50000);
     }
 
@@ -48,10 +63,11 @@ class OppoHotStoreTest {
         assertFalse(hot.get("mediaParams").has("unused"));
         assertCallbackFields(hot);
 
-        ArgumentCaptor<MapRecord<String, Object, Object>> event = ArgumentCaptor.forClass(MapRecord.class);
-        verify(stream).add(event.capture());
-        assertEquals("bid", event.getValue().getValue().get("type"));
-        JsonNode full = mapper.readTree((String) event.getValue().getValue().get("payload"));
+        ArgumentCaptor<MapRecord<byte[], byte[], byte[]>> event = ArgumentCaptor.forClass(MapRecord.class);
+        verify(streamCommands).xAdd(event.capture());
+        Map<byte[], byte[]> eventBody = event.getValue().getValue();
+        assertEquals("bid", new String(eventValue(eventBody, "type")));
+        JsonNode full = mapper.readTree(new String(eventValue(eventBody, "payload")));
         assertEquals("show-url", full.get("showUrls").get(0).asText());
         assertEquals("large-unused-value", full.get("mediaParams").get("unused").asText());
         assertEquals(123, full.get("price").asInt());
@@ -97,15 +113,26 @@ class OppoHotStoreTest {
     @Test
     void conversionLookupStillFallsBackToDb() {
         OppoBidRecord record = record();
+        when(values.multiGet(List.of("adx:oppo:bid:qk"))).thenReturn(java.util.Collections.singletonList(null));
         when(cold.getBidsByQks(List.of("qk"))).thenReturn(List.of(record));
         assertSame(record, store(Duration.ofHours(36)).findBidsByQks(List.of("qk")).get("qk"));
     }
 
     private JsonNode cachePayload(Duration ttl) throws Exception {
-        ArgumentCaptor<String> payload = ArgumentCaptor.forClass(String.class);
-        verify(values).set(eq("adx:oppo:bid:qk"), payload.capture(), eq(ttl));
-        verify(values).set(eq("adx:oppo:media:oppo:req"), eq(payload.getValue()), eq(ttl));
-        return mapper.readTree(payload.getValue());
+        ArgumentCaptor<byte[]> payload = ArgumentCaptor.forClass(byte[].class);
+        ArgumentCaptor<Expiration> expiration = ArgumentCaptor.forClass(Expiration.class);
+        verify(stringCommands).set(eq("adx:oppo:bid:qk".getBytes()), payload.capture(), expiration.capture(), eq(RedisStringCommands.SetOption.UPSERT));
+        verify(stringCommands).set(eq("adx:oppo:media:oppo:req".getBytes()), eq(payload.getValue()), any(Expiration.class), eq(RedisStringCommands.SetOption.UPSERT));
+        assertEquals(ttl.toMillis(), expiration.getValue().getExpirationTimeInMilliseconds());
+        return mapper.readTree(new String(payload.getValue()));
+    }
+
+    private byte[] eventValue(Map<byte[], byte[]> body, String key) {
+        for (Map.Entry<byte[], byte[]> entry : body.entrySet()) {
+            if (key.equals(new String(entry.getKey()))) return entry.getValue();
+        }
+        fail("missing stream field: " + key);
+        return null;
     }
 
     private void assertCallbackFields(JsonNode record) {