flink:通过Sink把数据写入kafka
·
package cn.edu.tju.demo;
import org.apache.flink.api.common.functions.*;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.collector.selector.OutputSelector;
import org.apache.flink.streaming.api.datastream.*;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.co.CoMapFunction;
import org.apache.flink.streaming.api.functions.source.SourceFunction;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer;
import org.apache.flink.util.Collector;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.*;
public class Test13 {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment environment = StreamExecutionEnvironment
.getExecutionEnvironment();
/* Properties properties = new Properties();
properties.setProperty(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
properties.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
properties.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"xx.xx.xx.xx:9092");*/
DataStreamSource<String> mySource = environment.addSource(new MySourceFunction());
mySource.addSink(new FlinkKafkaProducer<String>(
"xx.xx.xx.xx:9092","myTopic",new SimpleStringSchema( )));
environment.execute("my job");
}
public static class MySourceFunction implements SourceFunction<String> {
private boolean runningFlag = true;
@Override
public void run(SourceContext<String> ctx) throws Exception {
while (runningFlag){
ctx.collect("hi world");
ctx.collect("hello world");
Thread.sleep(30000);
}
}
@Override
public void cancel() {
runningFlag = false;
}
}
}
更多推荐
所有评论(0)