|
@@ -24,7 +24,7 @@ public class Main {
|
|
consumer.subscribe(List.of("article_843ea558d5ef41a1877584c62762632d_app_35","material","article_843ea558d5ef41a1877584c62762632d_app"));
|
|
consumer.subscribe(List.of("article_843ea558d5ef41a1877584c62762632d_app_35","material","article_843ea558d5ef41a1877584c62762632d_app"));
|
|
while (true) {
|
|
while (true) {
|
|
var kafkaDataList = new ArrayList<KafkaData>();
|
|
var kafkaDataList = new ArrayList<KafkaData>();
|
|
- for (ConsumerRecord<String, String> record : consumer.poll(Duration.ofMillis(2000))) {
|
|
|
|
|
|
+ for (ConsumerRecord<String, String> record : consumer.poll(Duration.ofMillis(60000))) {
|
|
kafkaDataList.add(new KafkaData(record.topic(),record.offset(), record.value(), new Timestamp(record.timestamp())));
|
|
kafkaDataList.add(new KafkaData(record.topic(),record.offset(), record.value(), new Timestamp(record.timestamp())));
|
|
}
|
|
}
|
|
for (Map.Entry<String, List<KafkaData>> entry : CollStreamUtil.groupByKey(kafkaDataList, KafkaData::getTopic).entrySet()) {
|
|
for (Map.Entry<String, List<KafkaData>> entry : CollStreamUtil.groupByKey(kafkaDataList, KafkaData::getTopic).entrySet()) {
|