Jak kodować / dekodować wiadomości Kafka za pomocą Avro binary encoder?
Próbuję użyć Avro do odczytywania/pisania wiadomości do Kafki. Ma ktoś może przykład użycia kodera binarnego Avro do kodowania/dekodowania danych, które będą umieszczane w kolejce komunikatów?
Potrzebuję części Avro bardziej niż części Kafki. A może powinienem spojrzeć na inne rozwiązanie? Zasadniczo staram się znaleźć bardziej efektywne rozwiązanie dla JSON w odniesieniu do przestrzeni. Avro właśnie wspomniano, ponieważ może być bardziej kompaktowy niż JSON.
5 answers
To jest podstawowy przykład. Nie próbowałem z wieloma partycjami / tematami.
//przykładowy kod producenta
import org.apache.avro.Schema;
import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.GenericRecord;
import org.apache.avro.io.*;
import org.apache.avro.specific.SpecificDatumReader;
import org.apache.avro.specific.SpecificDatumWriter;
import org.apache.commons.codec.DecoderException;
import org.apache.commons.codec.binary.Hex;
import kafka.javaapi.producer.Producer;
import kafka.producer.KeyedMessage;
import kafka.producer.ProducerConfig;
import java.io.ByteArrayOutputStream;
import java.io.File;
import java.io.IOException;
import java.nio.charset.Charset;
import java.util.Properties;
public class ProducerTest {
void producer(Schema schema) throws IOException {
Properties props = new Properties();
props.put("metadata.broker.list", "0:9092");
props.put("serializer.class", "kafka.serializer.DefaultEncoder");
props.put("request.required.acks", "1");
ProducerConfig config = new ProducerConfig(props);
Producer<String, byte[]> producer = new Producer<String, byte[]>(config);
GenericRecord payload1 = new GenericData.Record(schema);
//Step2 : Put data in that genericrecord object
payload1.put("desc", "'testdata'");
//payload1.put("name", "अasa");
payload1.put("name", "dbevent1");
payload1.put("id", 111);
System.out.println("Original Message : "+ payload1);
//Step3 : Serialize the object to a bytearray
DatumWriter<GenericRecord>writer = new SpecificDatumWriter<GenericRecord>(schema);
ByteArrayOutputStream out = new ByteArrayOutputStream();
BinaryEncoder encoder = EncoderFactory.get().binaryEncoder(out, null);
writer.write(payload1, encoder);
encoder.flush();
out.close();
byte[] serializedBytes = out.toByteArray();
System.out.println("Sending message in bytes : " + serializedBytes);
//String serializedHex = Hex.encodeHexString(serializedBytes);
//System.out.println("Serialized Hex String : " + serializedHex);
KeyedMessage<String, byte[]> message = new KeyedMessage<String, byte[]>("page_views", serializedBytes);
producer.send(message);
producer.close();
}
public static void main(String[] args) throws IOException, DecoderException {
ProducerTest test = new ProducerTest();
Schema schema = new Schema.Parser().parse(new File("src/test_schema.avsc"));
test.producer(schema);
}
}
//przykładowy kod konsumenta
Część 1: Kod grupy konsumenckiej: ponieważ możesz mieć więcej niż wielu konsumentów dla wielu partycji/ tematów.
import kafka.consumer.ConsumerConfig;
import kafka.consumer.KafkaStream;
import kafka.javaapi.consumer.ConsumerConnector;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Properties;
import java.util.concurrent.Executor;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
/**
* Created by on 9/1/15.
*/
public class ConsumerGroupExample {
private final ConsumerConnector consumer;
private final String topic;
private ExecutorService executor;
public ConsumerGroupExample(String a_zookeeper, String a_groupId, String a_topic){
consumer = kafka.consumer.Consumer.createJavaConsumerConnector(
createConsumerConfig(a_zookeeper, a_groupId));
this.topic = a_topic;
}
private static ConsumerConfig createConsumerConfig(String a_zookeeper, String a_groupId){
Properties props = new Properties();
props.put("zookeeper.connect", a_zookeeper);
props.put("group.id", a_groupId);
props.put("zookeeper.session.timeout.ms", "400");
props.put("zookeeper.sync.time.ms", "200");
props.put("auto.commit.interval.ms", "1000");
return new ConsumerConfig(props);
}
public void shutdown(){
if (consumer!=null) consumer.shutdown();
if (executor!=null) executor.shutdown();
System.out.println("Timed out waiting for consumer threads to shut down, exiting uncleanly");
try{
if(!executor.awaitTermination(5000, TimeUnit.MILLISECONDS)){
}
}catch(InterruptedException e){
System.out.println("Interrupted");
}
}
public void run(int a_numThreads){
//Make a map of topic as key and no. of threads for that topic
Map<String, Integer> topicCountMap = new HashMap<String, Integer>();
topicCountMap.put(topic, new Integer(a_numThreads));
//Create message streams for each topic
Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = consumer.createMessageStreams(topicCountMap);
List<KafkaStream<byte[], byte[]>> streams = consumerMap.get(topic);
//initialize thread pool
executor = Executors.newFixedThreadPool(a_numThreads);
//start consuming from thread
int threadNumber = 0;
for (final KafkaStream stream : streams) {
executor.submit(new ConsumerTest(stream, threadNumber));
threadNumber++;
}
}
public static void main(String[] args) {
String zooKeeper = args[0];
String groupId = args[1];
String topic = args[2];
int threads = Integer.parseInt(args[3]);
ConsumerGroupExample example = new ConsumerGroupExample(zooKeeper, groupId, topic);
example.run(threads);
try {
Thread.sleep(10000);
} catch (InterruptedException ie) {
}
example.shutdown();
}
}
Część 2: indywidualny konsument, który faktycznie konsumuje wiadomości.
import kafka.consumer.ConsumerIterator;
import kafka.consumer.KafkaStream;
import kafka.message.MessageAndMetadata;
import org.apache.avro.Schema;
import org.apache.avro.generic.GenericRecord;
import org.apache.avro.generic.IndexedRecord;
import org.apache.avro.io.DatumReader;
import org.apache.avro.io.Decoder;
import org.apache.avro.io.DecoderFactory;
import org.apache.avro.specific.SpecificDatumReader;
import org.apache.commons.codec.binary.Hex;
import java.io.File;
import java.io.IOException;
public class ConsumerTest implements Runnable{
private KafkaStream m_stream;
private int m_threadNumber;
public ConsumerTest(KafkaStream a_stream, int a_threadNumber) {
m_threadNumber = a_threadNumber;
m_stream = a_stream;
}
public void run(){
ConsumerIterator<byte[], byte[]>it = m_stream.iterator();
while(it.hasNext())
{
try {
//System.out.println("Encoded Message received : " + message_received);
//byte[] input = Hex.decodeHex(it.next().message().toString().toCharArray());
//System.out.println("Deserializied Byte array : " + input);
byte[] received_message = it.next().message();
System.out.println(received_message);
Schema schema = null;
schema = new Schema.Parser().parse(new File("src/test_schema.avsc"));
DatumReader<GenericRecord> reader = new SpecificDatumReader<GenericRecord>(schema);
Decoder decoder = DecoderFactory.get().binaryDecoder(received_message, null);
GenericRecord payload2 = null;
payload2 = reader.read(null, decoder);
System.out.println("Message received : " + payload2);
}catch (Exception e) {
e.printStackTrace();
System.out.println(e);
}
}
}
}
Test AVRO schema:
{
"namespace": "xyz.test",
"type": "record",
"name": "payload",
"fields":[
{
"name": "name", "type": "string"
},
{
"name": "id", "type": ["int", "null"]
},
{
"name": "desc", "type": ["string", "null"]
}
]
}
Ważne rzeczy do zapamiętania to:
Będziesz potrzebował standardu kafka i avro słoiki, aby uruchomić ten kod z pudełka.
Jest bardzo ważne rekwizyty.put ("serializer.Klasa", " kafka.serializer.DefaultEncoder"); Don
t use stringEncoder as that won
t działa, jeśli wysyłasz tablicę bajtów jako wiadomość.Można przekonwertować bajt [] na ciąg szesnastkowy i wysłać go, a na consumer reconvert ciąg szesnastkowy do bajtu [], a następnie do oryginalnej wiadomości.
Uruchom zookeeper i broker, jak wspomniano tutaj :- http://kafka.apache.org/documentation.html#quickstart i utworzyć temat o nazwie "page_views" lub cokolwiek chcesz.
Uruchom producenta.java, a następnie ConsumerGroupExample.java i zobaczyć dane avro są produkowane i zużywane.
Warning: date(): Invalid date.timezone value 'Europe/Kyiv', we selected the timezone 'UTC' for now. in /var/www/agent_stack/data/www/doraprojects.net/template/agent.layouts/content.php on line 54
2016-03-08 18:28:44
W końcu przypomniałem sobie, że zapytałem o listę dyskusyjną Kafki i otrzymałem następującą odpowiedź, która zadziałała idealnie.
Tak, możesz wysyłać wiadomości jako tablice bajtowe. Jeśli spojrzysz na konstruktor klasy Message, zobaczysz -
Def this (bytes: Array [Byte])
Teraz, patrząc na API producenta send ()-
Def send (producerData: ProducerData[K,V]*)
Możesz ustawić V jako typ wiadomości i K na to, co chcesz, aby twój klucz be. Jeśli nie zależy ci na partycjonowaniu za pomocą klucza, ustaw go na Message Typ też.
Dzięki, Neha
Warning: date(): Invalid date.timezone value 'Europe/Kyiv', we selected the timezone 'UTC' for now. in /var/www/agent_stack/data/www/doraprojects.net/template/agent.layouts/content.php on line 54
2011-12-01 21:03:07
Jeśli chcesz uzyskać tablicę bajtów z wiadomości Avro( część kafka jest już odebrana), użyj kodera binarnego:
GenericDatumWriter<GenericRecord> writer = new GenericDatumWriter<GenericRecord>(schema);
ByteArrayOutputStream os = new ByteArrayOutputStream();
try {
Encoder e = EncoderFactory.get().binaryEncoder(os, null);
writer.write(record, e);
e.flush();
byte[] byteData = os.toByteArray();
} finally {
os.close();
}
Warning: date(): Invalid date.timezone value 'Europe/Kyiv', we selected the timezone 'UTC' for now. in /var/www/agent_stack/data/www/doraprojects.net/template/agent.layouts/content.php on line 54
2014-07-22 01:48:31
Zaktualizowana Odpowiedź.
Kafka ma Avro serializer / deserializer ze współrzędnymi Mavena (sformatowanymi SBT):
"io.confluent" % "kafka-avro-serializer" % "3.0.0"
Przekazujesz instancję KafkaAvroSerializer do konstruktora KafkaProducer.
Następnie możesz utworzyć instancje Avro GenericRecord i używać ich jako wartości wewnątrz instancji Kafka ProducerRecord, które możesz wysłać za pomocą KafkaProducer.
Po stronie konsumentów Kafka, używasz KafkaAvroDeserializer i KafkaConsumer.
Warning: date(): Invalid date.timezone value 'Europe/Kyiv', we selected the timezone 'UTC' for now. in /var/www/agent_stack/data/www/doraprojects.net/template/agent.layouts/content.php on line 54
2016-06-09 03:55:47
Zamiast Avro, można również po prostu rozważyć kompresję danych; albo z gzip (dobra kompresja, wyższy procesor) lub LZF lub Snappy (znacznie szybsza, nieco wolniejsza kompresja).
Lub alternatywnie istnieje również Smile binary JSON, wspierany w Javie przez Jacksona (z to rozszerzenie): jest to kompaktowy format binarny i znacznie łatwiejszy w użyciu niż Avro:
ObjectMapper mapper = new ObjectMapper(new SmileFactory());
byte[] serialized = mapper.writeValueAsBytes(pojo);
// or back
SomeType pojo = mapper.readValue(serialized, SomeType.class);
W zasadzie ten sam kod co w JSON, z wyjątkiem przekazania innego formatu factory. Z perspektywy rozmiaru danych, to, czy Smile czy Avro są bardziej kompaktowe, zależy od szczegółów zastosowania; ale oba są bardziej kompaktowe niż JSON.
Zaletą jest to, że działa to szybko zarówno z JSON i Smile, z tym samym kodem, za pomocą tylko POJOs. W porównaniu do Avro, które wymaga generowania kodu, lub dużo ręcznego kodu do spakowania i rozpakowania GenericRecord
s.
Warning: date(): Invalid date.timezone value 'Europe/Kyiv', we selected the timezone 'UTC' for now. in /var/www/agent_stack/data/www/doraprojects.net/template/agent.layouts/content.php on line 54
2012-04-25 20:45:12