- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
我是 Spark 结构化流处理的新手,目前正在研究一个用例,其中结构化流应用程序将从 Azure IoT 中心-事件中心获取事件(例如每 20 秒后)。
任务是消耗这些事件并实时处理它。为此,我在 Spark-Java 中编写了下面的 Spark 结构化流程序。
以下是要点
问题:
...
public class EventSubscriber {
public static void main(String args[]) throws InterruptedException, StreamingQueryException {
String eventHubCompatibleEndpoint = "<My-EVENT HUB END POINT CONNECTION STRING>";
String connString = new ConnectionStringBuilder(eventHubCompatibleEndpoint).build();
EventHubsConf eventHubsConf = new EventHubsConf(connString).setConsumerGroup("$Default")
.setStartingPosition(EventPosition.fromEndOfStream()).setMaxRatePerPartition(100)
.setReceiverTimeout(java.time.Duration.ofMinutes(10));
SparkConf sparkConf = new SparkConf().setMaster("local[2]").setAppName("IoT Spark Streaming");
SparkSession spSession = SparkSession.builder()
.appName("IoT Spark Streaming")
.config(sparkConf).config("spark.sql.streaming.checkpointLocation", "<MY-CHECKPOINT-LOCATION>")
.getOrCreate();
Dataset<Row> inputStreamDF = spSession.readStream()
.format("eventhubs")
.options(eventHubsConf.toMap())
.load();
Dataset<Row> bodyRow = inputStreamDF.withColumn("body", new Column("body").cast(DataTypes.StringType)).select("body");
StructType jsonStruct = new StructType()
.add("eventType", DataTypes.StringType)
.add("payload", DataTypes.StringType);
Dataset<Row> messageRow = bodyRow.map((MapFunction<Row, Row>) value -> {
String valStr = value.getString(0).toString();
String payload = valStr;
Gson gson = new GsonBuilder().serializeNulls().setPrettyPrinting().create();
JsonObject jsonObj = gson.fromJson(valStr, JsonObject.class);
JsonElement methodName = jsonObj.get("method");
String eventType = null;
if(methodName != null) {
eventType = "OTHER_EVENT";
} else {
eventType = "DEVICE_EVENT";
}
Row jsonRow = RowFactory.create(eventType, payload);
return jsonRow;
}, RowEncoder.apply(jsonStruct));
messageRow.printSchema();
Dataset<Row> deviceEventRowDS = messageRow.filter("eventType = 'DEVICE_EVENT'");
deviceEventRowDS.printSchema();
Dataset<DeviceEvent> deviceEventDS = deviceEventRowDS.map((MapFunction<Row, DeviceEvent>) value -> {
String jsonString = value.getString(1).toString();
Gson gson = new GsonBuilder().serializeNulls().setPrettyPrinting().create();
DeviceMessage deviceMessage = gson.fromJson(jsonString, DeviceMessage.class);
DeviceEvent deviceEvent = deviceMessage.getDeviceEvent();
return deviceEvent;
}, Encoders.bean(DeviceEvent.class));
deviceEventDS.printSchema();
Dataset<Row> messageDataset = deviceEventDS.select(
functions.col("eventType"),
functions.col("deviceID"),
functions.col("description"),
functions.to_timestamp(functions.col("eventDate"), "yyyy-MM-dd hh:mm:ss").as("eventDate"),
functions.col("deviceModel"),
functions.col("pingRate"))
.select("eventType", "deviceID", "description", "eventDate", "deviceModel", "pingRate");
messageDataset.printSchema();
Dataset<Row> devWindowDataset = messageDataset.withWatermark("eventDate", "10 minutes")
.groupBy(functions.col("deviceID"),
functions.window(
functions.col("eventDate"), "10 minutes", "5 minutes"))
.count();
devWindowDataset.printSchema();
StreamingQuery query = devWindowDataset.writeStream().outputMode("append")
.format("parquet")
.option("truncate", "false")
.option("path", "<MY-PARQUET-FILE-OUTPUT-LOCATION>")
.start();
query.awaitTermination();
}}
...
与此相关的任何帮助或指示都会有用。
感谢和问候,
阿维纳什·德什穆克
最佳答案
Is it possible to store multiple events in parquet format in a file based on the window time?
是的。
How does the window operation works in this case?
以下代码是 Spark Structured Streaming 应用程序的主要部分:
Dataset<Row> devWindowDataset = messageDataset
.withWatermark("eventDate", "10 minutes")
.groupBy(
functions.col("deviceID"),
functions.window(functions.col("eventDate"), "10 minutes", "5 minutes"))
.count();
这表示底层状态存储应将每个 deviceID
和 eventDate
的状态保留 10 分钟以及额外的 10 分钟
(每个 withWatermark
)对于迟到的事件。换句话说,一旦事件的 eventDate
比流式查询开始时间晚 20 分钟,您就应该看到结果。
withWatermark
适用于后期事件,因此即使 groupBy
生成结果,只有在超过水印阈值后才会发出结果。
每 10 分钟(+ 10 分钟水印)应用相同的过程,并使用 5 分钟的窗口幻灯片。
将带有 window
运算符的 groupBy
视为多列聚合。
Also I would like to check the event state with previous event and based on some calculations (say event is not received by 5 minutes) I want to update the state.
这听起来像是 KeyValueGroupedDataset.flatMapGroupsWithState 的一个用例运算符(又名任意状态流聚合)。咨询Arbitrary Stateful Operations .
您也可能只想要众多 aggregation standard functions 之一。或 user-defined aggregation function (UDAF) .
关于java - Spark Structured Streaming - 有状态流处理中使用窗口操作进行事件处理,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/58444825/
https://github.com/mattdiamond/Recorderjs/blob/master/recorder.js中的代码 我不明白 JavaScript 语法,比如 (functio
在 iOS 7 及更早版本中,如果我们想在应用程序中找到 topMostWindow,我们通常使用以下代码行 [[[UIApplication sharedApplication] windows]
我已经尝试解决这个问题很长一段时间了:我无法访问窗口的 url,因为它位于另一个域上..有一些解决方案吗? function login() { var cb = window.ope
是否可以将 FFMPEG 视频流传递到 C# 窗口?现在它在新窗口中作为新进程打开,我只是想将它传递给我自己的 SessionWindow。 此时我像这样执行ffplay: public void E
我有一个名为 x 的矩阵看起来像这样: pTime Close 1 1275087600 1.2268 2 1275264000 1.2264 3 1275264300 1.2
在编译时,发生搜索,grep搜索等,Emacs会在单独的窗口中创建一个新的缓冲区来显示结果,有没有自动跳转到那个窗口的方法?这很有用,因为我可以使用 n 和 p 而不是 M-g n 和 M-g p 移
我有一个启动 PowerShell 脚本的批处理文件。 批处理文件: START Powershell -executionpolicy RemoteSigned -noexit -file "MyS
我有一个基于菜单栏的应用程序,单击图标时会显示一个窗口。在 Mac OS X Lion 上一切正常,但由于某种原因,在 Snow Leopard 和早期版本的 Mac OS X 上会出现错误。任何时候
在 macOS 中,如何在 Xcode 和/或 Interface Builder 中创建带有“集成标题栏和工具栏”的窗口? 这是“宽标题栏”类型的窗口,已添加到 OS X 10.10 Yosemit
在浏览器 (Chrome) 中 JavaScript: var DataModler = { Data: { Something: 'value' }, Process: functi
我有 3 个 html 页面。第 1 页链接到第 2 页,第 2 页链接到第 3 页(为了简单起见)。 我希望页面 2 中的链接打开页面 3 并关闭页面 1(选项卡 1)。 据我了解,您无法使用 Ja
当点击“创建节点”按钮时,如何打开一个新的框架或窗口?我希望新框架包含一个文本字段和下拉菜单,以便用户可以选择一个选项。 Create node Search node
我有一个用户控件,用于编辑应用程序中的某些对象。 我最近遇到一个实例,我想弹出一个新的对话框(窗口)来托管此用户控件。 如何实例化新窗口并将需要设置的任何属性从窗口传递到用户控件? 感谢您的宝贵时间。
我有一个Observable,它发出许多对象,我想使用window或buffer操作对这些对象进行分组。但是,我不想指定count参数来确定窗口中应包含多少个对象,而是希望能够使用自定义条件。 例如,
我有以下代码,它打开一个新的 JavaFX 阶段(我们称之为窗口)。 openAlertBox.setOnAction(e -> { AlertBox alert = AlertBox
我要添加一个“在新窗口中打开”上下文菜单项,该菜单项将以新的UIScene打开我的应用程序文档之一。当然,我只想在实际上支持多个场景的设备上显示该菜单项。 目前,我只是在检查设备是否是使用旧设备的iP
我正在尝试创建一个 AIR 应用程序来记录应用程序的使用情况,使用 AIR 从系统获取信息的唯一简单方法是使用命令行工具和抓取 标准输出 . 我知道像 这样的工具顶部 和 ps 对于 OS X,但它们
所以我有这个简单的 turtle 螺旋制作器,我想知道是否有一种方法可以打印出由该程序创建的我的设计副本。 代码: import turtle x= float(input("Angle: ")) y
我正在编写一个 C# WPF 程序,它将文本消息发送到另一个程序的窗口。我有一个宏程序作为我的键盘驱动程序 (Logitech g15) 的一部分,它已经这样做了,尽管它不会将击键直接发送到进程,而是
我尝试使用以下代码通过 UDP 发送,但得到了奇怪的结果。 if((sendto(newSocket, sendBuf, totalLength, 0, (SOCKADDR *)&sendAd
我是一名优秀的程序员,十分优秀!