代码之家  ›  专栏  ›  技术社区  ›  VladoDemcak

火花结构流动态串滤波器

  •  4
  • VladoDemcak  · 技术社区  · 8 年前

    我们正在尝试对结构化流应用程序使用动态过滤器。

    假设我们有以下Spark结构化流应用程序的伪实现:

    spark.readStream()
         .format("kafka")
         .option(...)
         ...
         .load()
         .filter(getFilter()) <-- dynamic staff - def filter(conditionExpr: String):
         .writeStream()
         .format("kafka")
         .option(.....)
         .start();
    

    getfilter返回字符串

    String getFilter() {
       // dynamic staff to create expression
       return expression; // eg. "column = true";
    }
    

    在当前版本的Spark中,是否可能存在动态过滤条件?我是说 getFilter() 方法应该动态返回一个过滤条件(假设它每10分钟刷新一次)。我们试图研究广播变量,但不确定结构化流是否支持这种情况。

    似乎提交作业后无法更新其配置。作为一种部署,我们使用 yarn .

    非常感谢您的每一个建议/选择。


    编辑: 假定 获取筛选器() 返回:

    (columnA = 1 AND columnB = true) OR customHiveUDF(columnC, 'input') != 'required' OR columnD > 8
    

    10分钟后,我们可以有小的变化(在第一个或之前没有第一个表达式),并且可能我们可以有一个新的表达式。( columnA = 2 )例如:

    customHiveUDF(columnC, 'input') != 'required' OR columnD > 10 OR columnA = 2
    

    目标是为一个Spark应用程序提供多个过滤器,并且不要提交多个作业。

    2 回复  |  直到 7 年前
        1
  •  1
  •   T. Gawęda    8 年前

    广播变量应该在这里正常。您可以编写如下类型的过滤器:

    query.filter(x => x > bv.value).writeStream(...)
    

    其中bv是a Broadcast 变量。您可以按照下面的描述进行更新: How can I update a broadcast variable in spark streaming?

    另一种解决方案是提供RCP或RESTful端点,并每隔10分钟询问该端点。例如(Java,因为这里更简单):

    class EndpointProxy {
    
         Configuration lastValue;
         long lastUpdated
         public static Configuration getConfiguration (){
    
              if (lastUpdated + refreshRate > System.currentTimeMillis()){
                   lastUpdated = System.currentTimeMillis();
                   lastValue = askMyAPI();
              }
              return lastValue;
         }
    }
    
    
    query.filter (x => x > EndpointProxy.getConfiguration().getX()).writeStream()
    

    编辑:针对用户问题的黑客解决方案:

    您可以创建一行视图: //confsdf应该在某个驱动程序端singleton中 var confsdf=seq(some content).todf(“somecolumn”)

    and then use:
    query.crossJoin(confsDF.as("conf")) // cross join as we have only 1 value 
          .filter("hiveUDF(conf.someColumn)")
          .writeStream()...
    
     new Thread() {
         confsDF = Seq(some new data).toDF("someColumn)
     }.start();
    

    这个黑客程序依赖于Spark默认的执行模型——微补丁。在每个触发器中,查询都将被重建,因此应该考虑新的数据。

    您也可以在线程中执行以下操作:

    Seq(some new data).toDF("someColumn).createOrReplaceTempView("conf")
    

    然后在查询中:

    .crossJoin(spark.table("conf"))
    

    两者都应该有效。记住,它不适用于连续处理模式

        2
  •  0
  •   Kaushal    8 年前

    下面是一个简单的例子,我在其中对来自套接字的记录进行动态过滤。代替日期,您可以使用任何RESTAPI来动态更新您的过滤器或轻量级ZooKeeper实例。

    注意 :-如果您计划使用任何RESTAPI或ZooKeeper或任何其他选项,请使用mapPartition而不是filter,因为在这种情况下,您已经为一个分区调用了一次API/连接。

    val lines = spark.readStream
      .format("socket")
      .option("host", "localhost")
      .option("port", 9999)
      .load()
    
    // Split the lines into words
    val words = lines.as[String].filter(_ == new java.util.Date().getMinutes.toString)
    
    // Generate running word count
    val wordCounts = words.groupBy("value").count()
    
    val query = wordCounts.writeStream
      .outputMode("complete")
      .format("console")
      .start()
    
    query.awaitTermination()