修改类型
This commit is contained in:
@@ -61,7 +61,7 @@ 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<Object> messageList,
|
public void listen(@Payload List<String> messageList,
|
||||||
@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
|
||||||
) {
|
) {
|
||||||
@@ -73,7 +73,7 @@ public class CorpusProcessKafkaProducer {
|
|||||||
}
|
}
|
||||||
try {
|
try {
|
||||||
for (int i = 0; i < messageList.size(); i++) {
|
for (int i = 0; i < messageList.size(); i++) {
|
||||||
String message =String.valueOf( messageList.get(i));
|
String message = messageList.get(i);
|
||||||
log.debug("corpusProcessKafkaProducer message:{}", message);
|
log.debug("corpusProcessKafkaProducer message:{}", message);
|
||||||
int partition = partitions.get(i);
|
int partition = partitions.get(i);
|
||||||
long offset = offsets.get(i);
|
long offset = offsets.get(i);
|
||||||
|
|||||||
Reference in New Issue
Block a user