liutong3 4 年 前
コミット
c5a1848e73

+ 7 - 3
comp-example/src/test/java/com/ai/ipu/example/kafka/KafkaCustomerExample.java

@ -14,11 +14,15 @@ public class KafkaCustomerExample {
14 14
15 15
    private static final int PULL_NUM = 100;
16 16
17
    private static final String SERVERS = "47.105.160.21:9091,47.105.160.21:9092,47.105.160.21:9093";
18
19
    private static final String DESERIALIZER = "org.apache.kafka.common.serialization.StringDeserializer";
20
17 21
    public static void main(String[] args) throws InterruptedException {
18 22
        Properties properties = new Properties();
19
        properties.put("bootstrap.servers", "47.105.160.21:9091,47.105.160.21:9092,47.105.160.21:9093");
20
        properties.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
21
        properties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
23
        properties.put("bootstrap.servers", SERVERS);
24
        properties.put("key.deserializer", DESERIALIZER);
25
        properties.put("value.deserializer", DESERIALIZER);
22 26
        properties.put(ConsumerConfig.GROUP_ID_CONFIG ,"test") ;
23 27
        properties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true");
24 28
        properties.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "1000");

+ 7 - 3
comp-example/src/test/java/com/ai/ipu/example/kafka/KafkaProducerExample.java

@ -13,11 +13,15 @@ import java.util.Scanner;
13 13
public class KafkaProducerExample {
14 14
    private static final Log logger = LogFactory.getLog(KafkaProducerExample.class);
15 15
16
    private static final String SERVERS = "47.105.160.21:9091,47.105.160.21:9092,47.105.160.21:9093";
17
18
    private static final String DESERIALIZER = "org.apache.kafka.common.serialization.StringDeserializer";
19
16 20
    public static void main(String[] args){
17 21
        Properties properties = new Properties();
18
        properties.put("bootstrap.servers", "47.105.160.21:9091,47.105.160.21:9092,47.105.160.21:9093");
19
        properties.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
20
        properties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
22
        properties.put("bootstrap.servers", SERVERS);
23
        properties.put("key.serializer", DESERIALIZER);
24
        properties.put("value.serializer", DESERIALIZER);
21 25
        properties.put("request.required.acks", "all");
22 26
        KafkaProducer<String, String> producer = new KafkaProducer<String, String>(properties);
23 27
        try {