修改日志
This commit is contained in:
@@ -73,7 +73,9 @@ public class TestController {
|
||||
public ResultMsg<Object> mockDccKafka(@RequestBody AicorpusTelephoneDTO message) {
|
||||
|
||||
for(int i=0;i<20;i++){
|
||||
message.setSourceId("000001"+i);
|
||||
String sourceId = "2000000"+i;
|
||||
message.setSourceId(sourceId);
|
||||
log.info("发送消息:"+sourceId);
|
||||
dccKafkaProducer.send("topic_voc_covert_text_log",JSONObject.toJSONString(message));
|
||||
}
|
||||
return ResultMsg.ok("ok");
|
||||
|
||||
@@ -17,8 +17,6 @@ import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
import org.springframework.cloud.context.config.annotation.RefreshScope;
|
||||
import org.springframework.kafka.annotation.KafkaListener;
|
||||
import org.springframework.kafka.support.KafkaHeaders;
|
||||
import org.springframework.messaging.handler.annotation.Header;
|
||||
import org.springframework.messaging.handler.annotation.Payload;
|
||||
import org.springframework.messaging.support.MessageBuilder;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
@@ -61,10 +59,7 @@ public class CorpusProcessKafkaProducer {
|
||||
|
||||
|
||||
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}")
|
||||
public void listen(@Payload List<String> messageList,
|
||||
@Header(KafkaHeaders.RECEIVED_PARTITION_ID) List<Integer> partitions,
|
||||
@Header(KafkaHeaders.OFFSET) List<Long> offsets
|
||||
) {
|
||||
public void listen(@Payload List<String> messageList) {
|
||||
long startTime = System.currentTimeMillis();
|
||||
log.info("corpusProcessKafkaProducerMessageSize:{}", messageList.size());
|
||||
if (CollectionUtils.isEmpty(messageList)) {
|
||||
@@ -72,13 +67,8 @@ public class CorpusProcessKafkaProducer {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
for (int i = 0; i < messageList.size(); i++) {
|
||||
String message = messageList.get(i);
|
||||
log.debug("corpusProcessKafkaProducer message:{}", message);
|
||||
int partition = partitions.get(i);
|
||||
long offset = offsets.get(i);
|
||||
|
||||
log.debug("corpusProcessKafkaProducer message: topic={}, partition={}, offset={}", "your-topic",partition, offset);
|
||||
for (String message:messageList) {
|
||||
log.debug("corpusProcessKafkaProducerMessage:{}", message);
|
||||
|
||||
// 提交异步任务
|
||||
CompletableFuture.runAsync(() -> {
|
||||
@@ -125,8 +115,7 @@ public class CorpusProcessKafkaProducer {
|
||||
* 处理 DCC 语料逻辑(异步发送 MQ)
|
||||
*/
|
||||
private void handleDccCorpus(TmTelephoneCorpus corpus, String rawMessage) {
|
||||
log.info("handleDccCorpus: categoryCode={}, sourceId={}",
|
||||
corpus.getCategoryCode(), corpus.getSourceId());
|
||||
log.info("corpusProcessKafkaProducerSaveSourceId: {}",corpus.getSourceId());
|
||||
|
||||
// 4. 保存语料
|
||||
tmTelephoneCorpusService.saveTelephoneCorpus(corpus);
|
||||
|
||||
Reference in New Issue
Block a user