Flink将数据写入到Kafka
·
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer;
import java.util.Properties;
/**
* Author:panghu
* Date:2022-05-29
* Description: 写数据到Kafka
*/
public class _16SinkToKafkaTest {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
DataStreamSource<String> streamSource = env.readTextFile("input/clicks.csv");
// kafka链接配置
Properties properties = new Properties();
properties.put("bootstrap.servers", "hadoop102:9092");
streamSource.addSink(new FlinkKafkaProducer<String>(
"clicks",
new SimpleStringSchema(),
properties
));
env.execute();
}
}
更多推荐



所有评论(0)