|
|
1
1
广播变量应该在这里正常。您可以编写如下类型的过滤器:
其中bv是a
另一种解决方案是提供RCP或RESTful端点,并每隔10分钟询问该端点。例如(Java,因为这里更简单):
编辑:针对用户问题的黑客解决方案: 您可以创建一行视图: //confsdf应该在某个驱动程序端singleton中 var confsdf=seq(some content).todf(“somecolumn”)
这个黑客程序依赖于Spark默认的执行模型——微补丁。在每个触发器中,查询都将被重建,因此应该考虑新的数据。 您也可以在线程中执行以下操作:
然后在查询中:
两者都应该有效。记住,它不适用于连续处理模式 |
|
|
2
0
下面是一个简单的例子,我在其中对来自套接字的记录进行动态过滤。代替日期,您可以使用任何RESTAPI来动态更新您的过滤器或轻量级ZooKeeper实例。 注意 :-如果您计划使用任何RESTAPI或ZooKeeper或任何其他选项,请使用mapPartition而不是filter,因为在这种情况下,您已经为一个分区调用了一次API/连接。
|