增加日志
This commit is contained in:
@@ -63,7 +63,6 @@ public class CorpusProcessKafkaProducer {
|
|||||||
|
|
||||||
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}")
|
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}")
|
||||||
public void listen( @Payload List<String> messageList,
|
public void listen( @Payload List<String> messageList,
|
||||||
@Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) List<String> keys,
|
|
||||||
@Header(KafkaHeaders.RECEIVED_PARTITION_ID) List<Integer> partitions,
|
@Header(KafkaHeaders.RECEIVED_PARTITION_ID) List<Integer> partitions,
|
||||||
@Header(KafkaHeaders.OFFSET) List<Long> offsets
|
@Header(KafkaHeaders.OFFSET) List<Long> offsets
|
||||||
) {
|
) {
|
||||||
@@ -77,11 +76,10 @@ public class CorpusProcessKafkaProducer {
|
|||||||
for (int i = 0; i < messageList.size(); i++) {
|
for (int i = 0; i < messageList.size(); i++) {
|
||||||
String message = messageList.get(i);
|
String message = messageList.get(i);
|
||||||
log.debug("corpusProcessKafkaProducer message:{}", message);
|
log.debug("corpusProcessKafkaProducer message:{}", message);
|
||||||
String key = keys.get(i);
|
|
||||||
int partition = partitions.get(i);
|
int partition = partitions.get(i);
|
||||||
long offset = offsets.get(i);
|
long offset = offsets.get(i);
|
||||||
|
|
||||||
log.debug("corpusProcessKafkaProducer message:key:{}, topic={}, partition={}, offset={}", "your-topic", key,partition, offset);
|
log.debug("corpusProcessKafkaProducer message: topic={}, partition={}, offset={}", "your-topic",partition, offset);
|
||||||
|
|
||||||
// 提交异步任务
|
// 提交异步任务
|
||||||
CompletableFuture.runAsync(() -> processSingleRecord(message), executor);
|
CompletableFuture.runAsync(() -> processSingleRecord(message), executor);
|
||||||
|
|||||||
Reference in New Issue
Block a user