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.
NotaO 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
-
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> -
Conecte-se ao mecanismo de streaming do Lindorm e grave dados. O exemplo completo de código está disponível abaixo:
NotaDados 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(); } } }