Flink测试kafka数据写入和打印出
  TEZNKK3IfmPf 2023年11月13日 22 0
package com.steven.flinkdemo.kafka;

import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer011;

import java.util.Properties;

public class KafkaConsumerApp {
public static void main(String[] args) throws Exception{
try {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
Properties properties = new Properties();
properties.setProperty("bootstrap.servers","dbos-bigdata-test006:9092,dbos-bigdata-test007:9092,dbos-bigdata-test008:9092");
properties.setProperty("group.id","flink");
DataStream<String> stream = env.addSource(new FlinkKafkaConsumer011<String>("test",new SimpleStringSchema(),properties));
stream.map(new MapFunction<String, String>() {
private static final long serialVersionUID = -6867736771747690202L;
@Override
public String map(String value) throws Exception {
return "flink:" + value;
}
}).print();
env.execute("consumer");
}catch (Exception e){
e.printStackTrace();
}
}
}
package com.steven.flinkdemo.kafka;

import org.apache.commons.lang.RandomStringUtils;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.source.SourceFunction;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer010;
import java.io.Serializable;
import java.util.Properties;

public class KafkaProdcerApp {
public static void main(String[] args) throws Exception{
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
Properties props = new Properties();
props.setProperty("bootstrap.servers","dbos-bigdata-test006:9092,dbos-bigdata-test007:9092,dbos-bigdata-test008:9092");
DataStream<String> stream = env.addSource(new SimpleStringGenerator());
stream.addSink(new FlinkKafkaProducer010<String>("test",new SimpleStringSchema(),props));
env.execute();

}
}
class SimpleStringGenerator implements SourceFunction<String>, Serializable{
private static final long serialVersionUID = 1L;
private volatile boolean isRunning = true;

@Override
public void run(SourceFunction.SourceContext<String> ctx) throws Exception{
while (isRunning){
String str = RandomStringUtils.randomAlphanumeric(5);
ctx.collect(str);
Thread.sleep(1000);
}
}
@Override
public void cancel(){
isRunning = false;
}
}
【版权声明】本文内容来自摩杜云社区用户原创、第三方投稿、转载,内容版权归原作者所有。本网站的目的在于传递更多信息,不拥有版权,亦不承担相应法律责任。如果您发现本社区中有涉嫌抄袭的内容,欢迎发送邮件进行举报,并提供相关证据,一经查实,本社区将立刻删除涉嫌侵权内容,举报邮箱: cloudbbs@moduyun.com

  1. 分享:
最后一次编辑于 2023年11月13日 0

暂无评论

推荐阅读
  TEZNKK3IfmPf   2023年11月15日   52   0   0 apachehadoopjava
  TEZNKK3IfmPf   2023年11月15日   31   0   0 apachejava
  TEZNKK3IfmPf   2023年11月15日   27   0   0 apachehadoop
TEZNKK3IfmPf