Todos os produtos
Search
Central de documentação

:Use um cliente Apache Kafka open source para gravar dados no mecanismo de streaming do Lindorm

Última atualização: Jun 28, 2026

A API do mecanismo de streaming do Lindorm é totalmente compatível com a API do Apache Kafka open source. Use a API do Apache Kafka para permitir que programas gravem dados no mecanismo de streaming do Lindorm. Ferramentas de terceiros open source, como Fluentd e Debezium, também coletam e gravam dados nesse mecanismo. Este tópico descreve como usar um cliente Apache Kafka open source para se conectar ao mecanismo de streaming do Lindorm e gravar dados, além de fornecer exemplos de código.

Pré-requisitos

  • Ambiente Java instalado com o Java Development Kit (JDK) 1,7 ou superior.

  • Endereço IP do cliente adicionado à lista de permissões da instância do Lindorm. Para mais informações, consulte Configure a whitelist.

  • Valor do Lindorm Stream Kafka Endpoint obtido. Para mais informações, consulte View endpoints.

    Nota

    O Lindorm Stream Kafka Endpoint especifica um endpoint de virtual private cloud (VPC) do mecanismo de streaming do Lindorm. Certifique-se de que a aplicação e a instância do Lindorm estejam implantadas na mesma VPC.

Procedimento

  1. Baixe um cliente Apache Kafka open source. Adicione as dependências do Maven ao arquivo pom.xml conforme o exemplo de código a seguir:

    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka-clients</artifactId>
        <version>0.10.2.2</version>
    </dependency>
  2. Conecte-se ao mecanismo de streaming do Lindorm e grave dados. O exemplo completo de código está disponível abaixo:

    Nota
    • Dados nos formatos JSON, Avro ou CSV podem ser gravados no mecanismo de streaming do Lindorm.

    • O valor do Lindorm Stream Kafka Endpoint no exemplo de código é um endpoint de VPC. Para obter informações sobre como obter o endpoint, consulte View endpoints.

    import org.apache.kafka.clients.producer.KafkaProducer;
    import org.apache.kafka.clients.producer.ProducerConfig;
    import org.apache.kafka.clients.producer.ProducerRecord;
    import org.apache.kafka.clients.producer.RecordMetadata;
    import org.codehaus.jettison.json.JSONObject;
    
    import java.util.Properties;
    import java.util.concurrent.Future;
    
    public class KafkaToLindormStreamDemo {
    
        public static void main(String[] args) {
            Properties props = new Properties();
    
            // Configure Lindorm Stream Kafka Endpoint. The value of Lindorm Stream Kafka Endpoint is a VPC endpoint. Make sure that your application and your Lindorm instance are deployed in the same VPC. 
            props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "Lindorm Stream Kafka Endpoint");
           // Specify the topic in which you want to store the physical data of your streaming data table.
            String topic = "log_topic";
    
            props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
            props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
            KafkaProducer<String, String> producer = new KafkaProducer<String, String>(props);
            try {
                JSONObject json = new JSONObject();
                // Write data to the streaming engine.
                json.put("timestamp", System.currentTimeMillis());
                json.put("loglevel", "ERROR");
                json.put("thread", "[ReportFinishedTask7-thread-4]");
                json.put("class", "engine.ImporterTaskManager(318)");
                json.put("detail", "Remove tasks fail: job name=e35318e5-52ea-48ab-ad2a-0144ffc6955e , task name=prepare_e35318e5-52ea-48ab-ad2a-0144ffc6955e , runningTasks=0");
                Future<RecordMetadata> future = producer.send(
              new ProducerRecord<String, String>(topic, json.getString("thread") + json.getLong("timestamp"),
                  json.toString()));
                producer.flush();
                try {
                    RecordMetadata recordMetadata = future.get();
                    System.out.println("Produce ok:" + recordMetadata.toString());
                } catch (Throwable t) {
                    System.out.println("Produce exception " + t.getMessage());
                    t.printStackTrace();
                }
            } catch (Exception e) {
                System.out.println("Produce exception " + e.getMessage());
                e.printStackTrace();
            }
        }
    }