【问题标题】:Flink: Left joining a stream with a static listFlink:左加入带有静态列表的流
【发布时间】:2020-11-26 23:15:37
【问题描述】:

我想将尝试流加入到被阻止电子邮件的静态列表中,并按 IP 对结果进行分组,以便以后可以计算一组相关统计信息。结果应以每 10 秒后 30 分钟的滑动窗口形式提供。以下是我尝试实现的几种方法之一:

override fun performQuery(): Table {
    val query = "SELECT ip, " +
        "COUNT(CASE WHEN success IS false THEN 1 END) AS fails, " +
        "COUNT(CASE WHEN success IS true THEN 1 END) AS successes, " +
        "COUNT(DISTINCT id) accounts, " +
        "COUNT(CASE WHEN id = 0 THEN 1 END) AS non_existing_accounts, " +
        "COUNT(CASE WHEN blockedEmail IS NOT NULL THEN 1 END) AS blocked_accounts " +
        "FROM Attempts " +
        "LEFT JOIN LATERAL TABLE(blockedEmailsList()) AS T(blockedEmail) ON TRUE " +
        "WHERE Attempts.email <> '' AND Attempts.createdAt < CURRENT_TIMESTAMP " +
        "GROUP BY HOP(Attempts.createdAt, INTERVAL '10' SECOND, INTERVAL '30' MINUTE), ip"

    return runQuery(query)
        .select("ip, accounts, fails, successes, non_existing_accounts, blocked_accounts")
}

这使用下面的用户定义表函数,它已经在我的tableEnv注册为blockedEmailsList

public class BlockedEmailsList extends TableFunction<Row> {
    private Collection<String> emails;

    public BlockedEmailsList(Collection<String> emails) {
        this.emails = emails;
    }

    public Row read(String email) {
        return Row.of(email);
    }

    public void eval() {
        this.emails.forEach(email -> collect(read(email)));
    }
}

但是,它返回以下错误:

Caused by: org.apache.flink.table.api.TableException: Rowtime attributes must not be in the input rows of a regular join. As a workaround you can cast the time attributes of input tables to TIMESTAMP before.

如果我按照建议进行操作并将created_at 转换为TIMESTAMP,我会得到这个:

org.apache.flink.client.program.ProgramInvocationException: The main method caused an error: Window can only be defined over a time attribute column.

我在 Stack Overflow 上发现了与这些异常相关的其他问题,但它们涉及流和时态表,并且它们都没有解决将流加入静态列表的情况。

有什么想法吗?

编辑:对于我的用例,Flink 项目中似乎有一个未解决的问题:https://cwiki.apache.org/confluence/display/FLINK/FLIP-17+Side+Inputs+for+DataStream+API

所以,我也接受解决方法建议。

【问题讨论】:

    标签: java kotlin apache-flink flink-streaming flink-sql


    【解决方案1】:
    Caused by: org.apache.flink.table.api.TableException: Rowtime attributes must not be in the input rows of a regular join. As a workaround you can cast the time attributes of input tables to TIMESTAMP before.
    

    原因是横向表函数是Flink正则连接,正则连接会发送空值,例如

    left:(K0, A), right(K1, T1)  => send    (K0, A, NULL, NULL)
    left:         , right(K0, T2) => retract (K0, A, NULL, NULL )  
                                       send   (K0, A, K0, T2)
    

    因此输入流中的时间属性在加入后会丢失。

    在您的情况下,您不需要 TableFunction,您可以使用标量函数 喜欢:

     public static class BlockedEmailFunction extends ScalarFunction {
         private static List<String> blockedEmails = ...;
         public Boolean eval(String email) {
            return blockedEmails.contains(attempt.getEmail());
         }
     }
    
    
    // register function
    env.createTemporarySystemFunction("blockedEmailFunction", BlockedEmailFunction.class);
    
    // call registered function in SQL and do window operation as your expected
    env.sqlQuery("SELECT blockedEmailFunction(email) as status, ip, createdAt FROM Attempts");
     
    

    【讨论】:

      【解决方案2】:

      我设法实施了解决我的问题的解决方法!

      我没有将流式尝试与电子邮件的静态列表连接起来,而是预先将每个尝试映射到一个带有 blockedEmail 属性的新尝试。如果静态列表blockedEmails 包含当前的尝试电子邮件,我将其blockedEmail 属性设置为true

      DataStream<Attempt> attemptsStream = sourceApi.<Attempt>startStream().map(new MapFunction<Attempt, Attempt>() {
          @Override
          public Attempt map(Attempt attempt) throws Exception {
              if (blockedEmails.contains(attempt.getEmail())) {
                  attempt.setBlockedEmail(true);
              }
              return attempt;
          }
      });
      

      静态列表blockedEmails 的类型为HashSet,因此查找时间为O(1)。

      最后,分组查询调整为:

      override fun performQuery(): Table {
          val query = "SELECT ip, " +
              "COUNT(CASE WHEN success IS false THEN 1 END) AS fails, " +
              "COUNT(CASE WHEN success IS true THEN 1 END) AS successes, " +
              "COUNT(DISTINCT id) accounts, " +
              "COUNT(CASE WHEN id = 0 THEN 1 END) AS non_existing_accounts, " +
              "COUNT(CASE WHEN blockedEmail IS true THEN 1 END) AS blocked_accounts " +
              "FROM Attempts " +
              "WHERE Attempts.email <> '' " +
              "GROUP BY HOP(Attempts.createdAt, INTERVAL '10' SECOND, INTERVAL '30' MINUTE), ip"
      
          return runQuery(query)
              .select("ip, accounts, fails, successes, non_existing_accounts, blocked_accounts")
      }
      

      到目前为止,流和静态列表之间的连接问题似乎尚未解决,但在我的情况下,上述解决方法解决了它。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2016-01-11
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2012-06-21
        • 1970-01-01
        • 2017-12-25
        相关资源
        最近更新 更多