【问题标题】:apache flink: how to interpret DataStream.print output?apache flink:如何解释 DataStream.print 输出?
【发布时间】:2015-11-27 03:54:32
【问题描述】:

我是 Flink 的新手,试图了解如何最有效地使用它。

我正在尝试使用 Window API,从 CSV 文件中读取。读取的行被转换为一个案例类,因此

case class IncomingDataUnit (
sensorUUID: String, radiationLevel: Int,photoSensor: Float,
humidity: Float,timeStamp: Long, ambientTemperature: Float)
  extends Serializable {

}

而且,这就是我读取行的方式:

env.readTextFile(inputPath).map(datum => {
      val fields = datum.split(",")
      IncomingDataUnit(
        fields(0),              // sensorUUID
        fields(1).toInt,        // radiationLevel
        fields(2).toFloat,      // photoSensor
        fields(3).toFloat,      // humidity
        fields(4).toLong,       // timeStamp
        fields(5).toFloat       // ambientTemperature
      )
    })

稍后,使用一个简单的窗口,我尝试打印最大 环境温度 ,因此:

env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime)

val readings =
      readIncomingReadings(env,"./sampleIOTTiny.csv")
      .map(e => (e.sensorUUID,e.ambientTemperature))
      .timeWindowAll(Time.of(5,TimeUnit.MILLISECONDS))
      .trigger(CountTrigger.of(5))
      .evictor(CountEvictor.of(4))
      .max(1)

readings.print

输出包含这些(来自一堆 DEBUG 日志语句):

1> (probe-987f2cb6,29.43)
1> (probe-987f2cb6,29.43)
3> (probe-dccefede,30.02)
3> (probe-42a9ddca,22.07)
2> (probe-df2d4cad,22.87)
2> (probe-20c609fb,27.62)
4> (probe-dccefede,30.02)

我想了解的是,人们如何解释这一点? repeated 1>s 代表什么?

让我感到困惑的是 probe-987f2cb6 在我的数据集中与 环境温度 29.43 不对应。它对应一个不同的值(准确地说是 14.72)。

仅供参考,以下是数据集:

probe-f076c2b0,201,842.53,75.5372,1448028160,29.37
probe-dccefede,199,749.25,78.6057,1448028160,27.46
probe-f29f9662,199,821.81,81.7831,1448028160,22.35
probe-5dac1d9f,195,870.71,83.1028,1448028160,15.98
probe-6c75cfbe,198,830.06,82.5607,1448028160,30.02
probe-4d78b545,204,778.42,78.412,1448028160,25.92
probe-400c5cdf,204,711.65,73.585,1448028160,22.18
probe-df2d4cad,199,820.8,72.936,1448028161,16.18
probe-f4ef109e,199,785.68,77.5647,1448028161,16.36
probe-3fac3350,200,720.12,78.2073,1448028161,19.19
probe-42a9ddca,193,819.12,74.3712,1448028161,22.07
probe-252a5bbd,197,710.32,80.6072,1448028161,14.64
probe-987f2cb6,200,750.4,76.0533,1448028161,14.72
probe-24444323,197,816.06,84.0816,1448028161,4.405
probe-6dd6fdc4,201,717.64,78.4031,1448028161,29.43
probe-20c609fb,204,804.37,84.5243,1448028161,22.87
probe-c027fdc9,195,858.61,81.7682,1448028161,24.47
probe-2c6cd3de,198,826.96,85.26,1448028162,18.99
probe-960906ca,197,797.63,77.4359,1448028162,27.62

我可能缺少很多东西。请刺激我。

忘了说:我使用的是Flink 0.10.0。

【问题讨论】:

    标签: scala apache-flink flink-streaming


    【解决方案1】:

    “>X”表示打印结果元组的并行任务的任务 ID。我只是想知道为什么输出显示值 1 到 4 - 因为您使用的是 non-parallel 窗口(数据流未通过 .keyBy() 分区),我希望打印是链式和非并行的。但也许不是,4 个并行打印任务正在运行。

    关于您的结果:如果窗口触发,则在所有元组上计算字段 1 的最大值,并将窗口的 head 元组值用于字段 0。如果要返回完整包含最大值的元组,您可以使用maxBy() 而不是max()

    【讨论】:

    • @matthias-j-sax:既然你提到了,我鼓起勇气说这也是我的问题 - 为什么在我没有特别要求任何并行性的情况下执行 4 个不同的任务。至于 maxBy 而不是 max,我会试试的。但我想知道将 maxed field '1' 与 head-reacord-field '0' 结合起来在语义上是否正确。这种解释有什么用?
    • Flink 总是根据你的硬件使用默认的并行度(我猜你有一台 4 核的机器)。尽管如此,非并行窗口将在单个线程中执行——但这不适用于您的其他运算符。您当然可以通过StreamExecutionEnvironment.setParallelism(int) 将默认并行度显式设置为一个(或为每个运算符单独设置:operator.setParallelism(int))。 max 在 key-streams 上很有用,可以输出聚合以及窗口中所有元组的 key-attribute。
    • 我们在这个问题中混合了两件事:对已回答的输出的解释和窗口语义。关于窗口语义:您正在使用处理时间,以及带有计数触发器和计数驱逐器的 5 毫秒时间窗口。通过指定计数触发器,您可以覆盖时间触发器,即窗口不会在 5 毫秒后触发。相反,它在 5 个元素到达后​​触发。在应用窗口函数并从窗口中删除前四个元素之前调用驱逐器。因此,max() 仅适用于最后收到的元素。
    • @FabianHueske:感谢您的解释。我现在更清楚了。当我接受答案时,我想这将归功于马蒂亚斯,但你们俩都花时间解释了。非常感谢。
    • 另外,我也发邮件给用户组,询问这个特定的窗口语义。我提到它只是为了让我的问题在那里看起来不重复。
    猜你喜欢
    • 1970-01-01
    • 2017-08-20
    • 1970-01-01
    • 2015-11-11
    • 2016-08-18
    • 2011-09-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多