flink使用kafka为数据源
1.10.2 org.apache.flink flink-java ${flink.version} org.apache.flink flink-streaming-java_2.12 ${flink.version} org.apache.flink flink-connector-kafka_2.12 ${flink.version} org.slf4j slf4j-log4j12 1.7.25 runtime log4j log4j 1.2.16 runtime
1.增加flink依赖
2.代码实现
import org.apache.flink.api.common.functions.MapFunction; import org.apache.flink.api.common.serialization.DeserializationSchema; import org.apache.flink.api.common.serialization.SimpleStringSchema; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.api.java.tuple.Tuple3; import org.apache.flink.shaded.guava18.com.google.common.base.Strings; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer; import java.util.Properties;
public class KafkaExample { public static void main(String[] args) { Properties properties = new Properties(); properties.setProperty("bootstrap.servers", "192.168.1.85:9092,192.168.1.86:9092,192.168.1.87:9092"); properties.setProperty("group.id", "g2"); DeserializationSchemadeserializationSchema = new SimpleStringSchema(); String topic = "customer_statusChangedEvent"; final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); DataStream text = env.addSource( new FlinkKafkaConsumer (topic, deserializationSchema, properties)); DataStream > dataStream = text .filter(p -> !Strings.isNullOrEmpty(p)) .map(new MapFunction >() { @Override public Tuple2 map(String p) throws Exception { CustomerStatusChangedEvent event = JsonHelper.fromJson(p, CustomerStatusChangedEvent.class); return new Tuple2 (Long.toString(event.getCustomerId()), event.getNewStatus()); } }); dataStream.print(); try { env.execute("Flink-Kafka"); } catch (Exception ex) { } } }
public class CustomerStatusChangedEvent { private Long customerId; @JsonProperty("nStatus") private Integer newStatus; @JsonProperty("oStatus") private Integer oldStatus; public Long getCustomerId() { return customerId; } public void setCustomerId(Long customerId) { this.customerId = customerId; } public Integer getNewStatus() { return newStatus; } public void setNewStatus(Integer newStatus) { this.newStatus = newStatus; } public Integer getOldStatus() { return oldStatus; } public void setOldStatus(Integer oldStatus) { this.oldStatus = oldStatus; } }
4.输出
当kafka的对应的topic有数据产生时, 将输出 类似下面数据
(5010,1)