-
Notifications
You must be signed in to change notification settings - Fork 17
/
KafkaFlinkAvroParquet.java
67 lines (53 loc) · 2.44 KB
/
KafkaFlinkAvroParquet.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
package it.okkam.flink;
import java.util.Arrays;
import java.util.Properties;
import org.apache.flink.api.common.serialization.DeserializationSchema;
import org.apache.flink.api.common.serialization.SerializationSchema;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.typeutils.TypeExtractor;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer010;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer010;
import org.okkam.flink.avro.AvroDeserializationSchema;
import org.okkam.flink.avro.AvroSerializationSchema;
import org.okkam.flink.parquet.avro.Person;
import org.okkam.flink.parquet.avro.PhoneNumber;
import org.okkam.flink.parquet.avro.PhoneType;
public class KafkaFlinkAvroParquet {
static String topicId = "test";
static String kafkaPort = "localhost:9092";
static String zkPort = "localhost:2181";
static SerializationSchema<Person> ser = new AvroSerializationSchema<>(Person.class);
static DeserializationSchema<Person> deser = new AvroDeserializationSchema<>(Person.class);
/**
* Kafka test.
*
* @param args
* the program arguments
*
* @throws Exception
* if any exception occurs
*/
public static void main(String[] args) throws Exception {
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
Person person = new Person();
person.setId(42);
person.setName("Felix");
person.setEmail("test@test.com");
person.setPhone(Arrays.asList(new PhoneNumber("123456", PhoneType.WORK)));
Properties properties = new Properties();
properties.setProperty("bootstrap.servers", kafkaPort);
properties.setProperty("group.id", topicId);
properties.setProperty("zookeeper.connect", zkPort);
properties.setProperty("batch.size", "0");
TypeInformation<Person> typeInfo = TypeExtractor.getForClass(Person.class);
DataStream<Person> stream = env.fromCollection(Arrays.asList(person), typeInfo);
stream.addSink(new FlinkKafkaProducer010<Person>(topicId, ser, properties));
DataStreamSource<Person> messageStream = env
.addSource(new FlinkKafkaConsumer010<Person>(topicId, deser, properties));
messageStream.print();
env.execute();
}
}