【发布时间】:2017-10-29 21:11:10
【问题描述】:
我正在尝试运行 Flink 流式传输作业。我想确定流处理的吞吐量和延迟。我已经启动了 Kafka 代理服务器并收到来自 kafka 的传入消息。我如何计算每秒消息数(吞吐量)? (比如rdd.count。有没有类似的方法来获取传入消息的数量)
(完整场景:我已通过 Producer 将消息作为 Json 对象发送。我正在添加一些信息,例如名称作为字符串以及 Json 对象中的 System.currentTimeMills。 流式传输时,如何通过messageStream(DataStream)获取发送的json对象?)
提前致谢。
代码:
/** * 从 Kafka 读取字符串并将它们打印到标准输出。 */
public static void main(String[] args) throws Exception {
System.setProperty("hadoop.home.dir", "c:/winutils/");
// parse input argum ents
final ParameterTool parameterTool = ParameterTool.fromArgs(args);
if(parameterTool.getNumberOfParameters() < 4) {
System.out.println("Missing parameters!\nUsage: Kafka --topic <topic> " +
"--bootstrap.servers <kafka brokers> --zookeeper.connect <zk quorum> --group.id <some id>");
return;
}
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.getConfig().disableSysoutLogging();
env.getConfig().setRestartStrategy(RestartStrategies.fixedDelayRestart(4, 10000));
env.enableCheckpointing(5000); // create a checkpoint every 5 seconds
env.getConfig().setGlobalJobParameters(parameterTool); // make parameters available in the web interface
DataStream<String> messageStream = env
.addSource(new FlinkKafkaConsumer010<>(
parameterTool.getRequired("topic"),
new SimpleStringSchema(),
parameterTool.getProperties()));
messageStream.print();
env.execute();
}
【问题讨论】:
标签: json apache-kafka apache-flink kafka-producer-api flink-streaming