【问题标题】:Flow time stamp through streaming functions通过流函数的流时间戳
【发布时间】:2017-03-07 09:30:24
【问题描述】:

每次使用 Spark Streaming 运行批处理时,如何/是否可以生成随机数或获取系统时间?

我有两个处理一批消息的函数: 1 - 首先处理密钥,创建一个文件(csv)并写入标题 2 - 秒处理每条消息并将数据添加到 csv

我希望将每个批次的文件存储在单独的文件夹中:

/output/folderBatch1/file1.csv, file2.csv, etc.csv
/output/folderBatch2/file1.csv, file2.csv, etc.csv
/output/folderBatch3/file1.csv, file2.csv, etc.csv

如何创建一个变量,甚至只是一个 Spark Streaming 可以使用的简单计数器?

下面的代码获取系统时间,但因为它是“普通 Java”,所以它只执行一次,并且在每次批处理运行时都是相同的值。

JavaPairInputDStream<String, byte[]> messages;
messages = KafkaUtils.createDirectStream(
        jssc,
        String.class,
        byte[].class,
        StringDecoder.class,
        DefaultDecoder.class,
        kafkaParams,
        topicsSet
);

/**
 * Declare what computation needs to be done
 */
JavaPairDStream<String, Iterable<byte[]>> groupedMessages = messages.groupByKey();

String time = Long.toString(System.currentTimeMillis());        //this is only ever run once and is the same value for each batch!

groupedMessages.map(new WriteHeaders(time)).print();

groupedMessages.map(new ProcessMessages(time)).print();

谢谢, 卡。

【问题讨论】:

  • 为什么不直接将System.currentTimeMillis 传递给groupedMessage.map(..)
  • 每个批次的两个地图函数的值需要相同。如果它不是相同的值,我将有带有标题的 csv 文件,没有数据和没有标题的数据文件。
  • 对于初学者,您可以考虑将这两个操作合并到一个 map 调用中。这样您就可以在本地共享时间戳。除此之外,您可以创建一个元组:Tuple2&lt;Long, Iterable&lt;byte[]&gt;&gt;,它会捕获时间戳。
  • 我试过了:Tuple2 timeStamp = new Tuple2(Long.toString(System.currentTimeMillis()), new byte[]{});但当然,这只会执行一次 - 我怎样才能以“Spark Streaming”的方式编写它?

标签: java apache-spark batch-processing spark-streaming


【解决方案1】:

您可以通过额外的map 调用来添加时间戳并将其传递。这意味着您将拥有Tuple2&lt;Long, Iterable&lt;byte[]&gt;) 的值,而不是Iterable&lt;byte[]&gt; 类型的值:

JavaDStream<Tuple2<String, Tuple2<Long, Iterable<byte[]>>>> groupedWithTimeStamp = 
  groupedMessages
    .map((Function<Tuple2<String, Iterable<byte[]>>, 
      Tuple2<String, Tuple2<Long, Iterable<byte[]>>>>) kvp -> 
        new Tuple2<>(kvp._1, new Tuple2<>(System.currentTimeMillis(), kvp._2)));

现在您在每个map 中都有时间戳,从现在开始,您可以通过以下方式访问它:

groupedWithTimeStamp.map(value -> value._2._1); // This will access the timestamp.

【讨论】:

    猜你喜欢
    • 2011-09-25
    • 1970-01-01
    • 1970-01-01
    • 2014-07-04
    • 1970-01-01
    • 2015-06-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多