Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
doujutun3207
flink
提交
cac9fa02
F
flink
项目概览
doujutun3207
/
flink
与 Fork 源项目一致
从无法访问的项目Fork
通知
24
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
F
flink
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
体验新版 GitCode,发现更多精彩内容 >>
提交
cac9fa02
编写于
3月 24, 2017
作者:
F
Fabian Hueske
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[hotfix] [table] Disable event-time OVER RANGE UNBOUNDED PRECEDING window.
上级
fe2c61a2
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
33 addition
and
28 deletion
+33
-28
flink-libraries/flink-table/src/main/scala/org/apache/flink/table/plan/nodes/datastream/DataStreamOverAggregate.scala
...table/plan/nodes/datastream/DataStreamOverAggregate.scala
+13
-8
flink-libraries/flink-table/src/test/scala/org/apache/flink/table/api/scala/stream/sql/SqlITCase.scala
...g/apache/flink/table/api/scala/stream/sql/SqlITCase.scala
+20
-20
未找到文件。
flink-libraries/flink-table/src/main/scala/org/apache/flink/table/plan/nodes/datastream/DataStreamOverAggregate.scala
浏览文件 @
cac9fa02
...
...
@@ -113,23 +113,28 @@ class DataStreamOverAggregate(
if
(
overWindow
.
isRows
)
{
// ROWS clause bounded OVER window
throw
new
TableException
(
"
ROWS clause bounded proc-time OVER window no
supported yet."
)
"
processing-time OVER ROWS PRECEDING window is not
supported yet."
)
}
else
{
// RANGE clause bounded OVER window
throw
new
TableException
(
"
RANGE clause bounded proc-time OVER window no
supported yet."
)
"
processing-time OVER RANGE PRECEDING window is not
supported yet."
)
}
}
else
{
throw
new
TableException
(
"OVER window only support ProcessingTime UNBOUNDED PRECEDING and CURRENT ROW "
+
"condition."
)
"processing-time OVER RANGE FOLLOWING window is not supported yet."
)
}
case
_:
RowTimeType
=>
// row-time OVER window
if
(
overWindow
.
lowerBound
.
isPreceding
&&
overWindow
.
lowerBound
.
isUnbounded
&&
overWindow
.
upperBound
.
isCurrentRow
)
{
// unbounded preceding OVER window
createUnboundedAndCurrentRowEventTimeOverWindow
(
inputDS
)
if
(
overWindow
.
isRows
)
{
// unbounded preceding OVER ROWS window
createUnboundedAndCurrentRowEventTimeOverWindow
(
inputDS
)
}
else
{
// unbounded preceding OVER RANGE window
throw
new
TableException
(
"row-time OVER RANGE UNBOUNDED PRECEDING window is not supported yet."
)
}
}
else
if
(
overWindow
.
lowerBound
.
isPreceding
&&
overWindow
.
upperBound
.
isCurrentRow
)
{
// bounded OVER window
if
(
overWindow
.
isRows
)
{
...
...
@@ -138,11 +143,11 @@ class DataStreamOverAggregate(
}
else
{
// RANGE clause bounded OVER window
throw
new
TableException
(
"
RANGE clause bounded row-time OVER window no
supported yet."
)
"
row-time OVER RANGE PRECEDING window is not
supported yet."
)
}
}
else
{
throw
new
TableException
(
"row-time OVER
window only support CURRENT ROW condition
."
)
"row-time OVER
RANGE FOLLOWING window is not supported yet
."
)
}
case
_
=>
throw
new
TableException
(
s
"Unsupported time type {$timeType}"
)
...
...
flink-libraries/flink-table/src/test/scala/org/apache/flink/table/api/scala/stream/sql/SqlITCase.scala
浏览文件 @
cac9fa02
...
...
@@ -448,15 +448,15 @@ class SqlITCase extends StreamingWithStateTestBase {
val
sqlQuery
=
"SELECT a, b, c, "
+
"SUM(b) over ("
+
"partition by a order by rowtime() r
ange
between unbounded preceding and current row), "
+
"partition by a order by rowtime() r
ows
between unbounded preceding and current row), "
+
"count(b) over ("
+
"partition by a order by rowtime() r
ange
between unbounded preceding and current row), "
+
"partition by a order by rowtime() r
ows
between unbounded preceding and current row), "
+
"avg(b) over ("
+
"partition by a order by rowtime() r
ange
between unbounded preceding and current row), "
+
"partition by a order by rowtime() r
ows
between unbounded preceding and current row), "
+
"max(b) over ("
+
"partition by a order by rowtime() r
ange
between unbounded preceding and current row), "
+
"partition by a order by rowtime() r
ows
between unbounded preceding and current row), "
+
"min(b) over ("
+
"partition by a order by rowtime() r
ange
between unbounded preceding and current row) "
+
"partition by a order by rowtime() r
ows
between unbounded preceding and current row) "
+
"from T1"
val
data
=
Seq
(
...
...
@@ -526,15 +526,15 @@ class SqlITCase extends StreamingWithStateTestBase {
val
sqlQuery
=
"SELECT a, b, c, "
+
"SUM(b) over ("
+
"partition by a order by rowtime() r
ange
between unbounded preceding and current row), "
+
"partition by a order by rowtime() r
ows
between unbounded preceding and current row), "
+
"count(b) over ("
+
"partition by a order by rowtime() r
ange
between unbounded preceding and current row), "
+
"partition by a order by rowtime() r
ows
between unbounded preceding and current row), "
+
"avg(b) over ("
+
"partition by a order by rowtime() r
ange
between unbounded preceding and current row), "
+
"partition by a order by rowtime() r
ows
between unbounded preceding and current row), "
+
"max(b) over ("
+
"partition by a order by rowtime() r
ange
between unbounded preceding and current row), "
+
"partition by a order by rowtime() r
ows
between unbounded preceding and current row), "
+
"min(b) over ("
+
"partition by a order by rowtime() r
ange
between unbounded preceding and current row) "
+
"partition by a order by rowtime() r
ows
between unbounded preceding and current row) "
+
"from T1"
val
data
=
Seq
(
...
...
@@ -596,11 +596,11 @@ class SqlITCase extends StreamingWithStateTestBase {
env
.
setParallelism
(
1
)
val
sqlQuery
=
"SELECT a, b, c, "
+
"SUM(b) over (order by rowtime() r
ange
between unbounded preceding and current row), "
+
"count(b) over (order by rowtime() r
ange
between unbounded preceding and current row), "
+
"avg(b) over (order by rowtime() r
ange
between unbounded preceding and current row), "
+
"max(b) over (order by rowtime() r
ange
between unbounded preceding and current row), "
+
"min(b) over (order by rowtime() r
ange
between unbounded preceding and current row) "
+
"SUM(b) over (order by rowtime() r
ows
between unbounded preceding and current row), "
+
"count(b) over (order by rowtime() r
ows
between unbounded preceding and current row), "
+
"avg(b) over (order by rowtime() r
ows
between unbounded preceding and current row), "
+
"max(b) over (order by rowtime() r
ows
between unbounded preceding and current row), "
+
"min(b) over (order by rowtime() r
ows
between unbounded preceding and current row) "
+
"from T1"
val
data
=
Seq
(
...
...
@@ -651,11 +651,11 @@ class SqlITCase extends StreamingWithStateTestBase {
env
.
setParallelism
(
1
)
val
sqlQuery
=
"SELECT a, b, c, "
+
"SUM(b) over (order by rowtime() r
ange
between unbounded preceding and current row), "
+
"count(b) over (order by rowtime() r
ange
between unbounded preceding and current row), "
+
"avg(b) over (order by rowtime() r
ange
between unbounded preceding and current row), "
+
"max(b) over (order by rowtime() r
ange
between unbounded preceding and current row), "
+
"min(b) over (order by rowtime() r
ange
between unbounded preceding and current row) "
+
"SUM(b) over (order by rowtime() r
ows
between unbounded preceding and current row), "
+
"count(b) over (order by rowtime() r
ows
between unbounded preceding and current row), "
+
"avg(b) over (order by rowtime() r
ows
between unbounded preceding and current row), "
+
"max(b) over (order by rowtime() r
ows
between unbounded preceding and current row), "
+
"min(b) over (order by rowtime() r
ows
between unbounded preceding and current row) "
+
"from T1"
val
data
=
Seq
(
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录