【发布时间】: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