【发布时间】:2019-04-17 10:40:25
【问题描述】:
我正在将流分析作业从 Databricks/Spark 迁移到 Azure 流分析。输入来自 IoTHub,每当传感器值在阈值范围之间变化时(例如,从“警告”范围变为“警报”范围),查询都必须发出事件。
现有解决方案利用“有状态流式传输”,即它在内存中保存每个设备的最后状态,并在每条新消息上进行比较。在作业启动时(或在某些其他情况下)没有“最后状态”;在这种情况下,会创建一个额外的事件 - 并由下游组件优雅地处理。
我正在尝试在 ASA 中实现此功能:
- 使用 可以轻松地与上一条记录进行比较
lag(value, 1, null) over (partition by(serialMachine) limit duration(minute, 60))
- 使用本地输入数据进行测试时,上述结果对于第一条记录为空,可用于创建消息。
- 但在 Azure 上运行时,“lag”会返回一个值,即使它的源记录的时间戳早于配置的作业开始时间。我猜它被视为“输出开始时间”,无论此时间戳如何,所有可用消息或至少更多消息都会从 IoTHub 加载。
我尝试了 ISFIRST 和 LAST 函数,但所有这些都指的是一个时间窗口,即派生条件将定期得到满足。但我只需要一次。
有什么解决方法的想法吗?
【问题讨论】: