Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
doujutun3207
flink
提交
0a66d09b
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,体验更适合开发者的 AI 搜索 >>
未验证
提交
0a66d09b
编写于
12月 18, 2020
作者:
A
Aljoscha Krettek
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[refactor] Rename StreamConfig.setTypeSerializersIn() to setupNetworkInputs()
Because that's what it does.
上级
28975540
变更
7
隐藏空白更改
内联
并排
Showing
7 changed file
with
11 addition
and
11 deletion
+11
-11
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamConfig.java
...va/org/apache/flink/streaming/api/graph/StreamConfig.java
+1
-1
flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/OneInputStreamTaskTest.java
...flink/streaming/runtime/tasks/OneInputStreamTaskTest.java
+1
-1
flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/OneInputStreamTaskTestHarness.java
...treaming/runtime/tasks/OneInputStreamTaskTestHarness.java
+1
-1
flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamConfigChainer.java
...he/flink/streaming/runtime/tasks/StreamConfigChainer.java
+1
-1
flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/TwoInputStreamTaskTestHarness.java
...treaming/runtime/tasks/TwoInputStreamTaskTestHarness.java
+1
-1
flink-streaming-java/src/test/java/org/apache/flink/streaming/util/OneInputStreamOperatorTestHarness.java
...ink/streaming/util/OneInputStreamOperatorTestHarness.java
+5
-5
flink-table/flink-table-runtime-blink/src/main/java/org/apache/flink/table/runtime/operators/multipleinput/MultipleInputStreamOperatorBase.java
...rators/multipleinput/MultipleInputStreamOperatorBase.java
+1
-1
未找到文件。
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamConfig.java
浏览文件 @
0a66d09b
...
...
@@ -248,7 +248,7 @@ public class StreamConfig implements Serializable {
}
}
public
void
set
TypeSerializersIn
(
TypeSerializer
<?>...
serializers
)
{
public
void
set
upNetworkInputs
(
TypeSerializer
<?>...
serializers
)
{
InputConfig
[]
inputs
=
new
InputConfig
[
serializers
.
length
];
for
(
int
i
=
0
;
i
<
serializers
.
length
;
i
++)
{
inputs
[
i
]
=
new
NetworkInputConfig
(
serializers
[
i
],
i
);
...
...
flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/OneInputStreamTaskTest.java
浏览文件 @
0a66d09b
...
...
@@ -997,7 +997,7 @@ public class OneInputStreamTaskTest extends TestLogger {
for
(
int
chainedIndex
=
1
;
chainedIndex
<
numberChainedTasks
;
chainedIndex
++)
{
TestingStreamOperator
<
Integer
,
Integer
>
chainedOperator
=
new
TestingStreamOperator
<>();
StreamConfig
chainedConfig
=
new
StreamConfig
(
new
Configuration
());
chainedConfig
.
set
TypeSerializersIn
(
StringSerializer
.
INSTANCE
);
chainedConfig
.
set
upNetworkInputs
(
StringSerializer
.
INSTANCE
);
chainedConfig
.
setStreamOperator
(
chainedOperator
);
chainedConfig
.
setOperatorID
(
new
OperatorID
(
0L
,
chainedIndex
));
chainedTaskConfigs
.
put
(
chainedIndex
,
chainedConfig
);
...
...
flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/OneInputStreamTaskTestHarness.java
浏览文件 @
0a66d09b
...
...
@@ -133,7 +133,7 @@ public class OneInputStreamTaskTestHarness<IN, OUT> extends StreamTaskTestHarnes
}
streamConfig
.
setNumberOfNetworkInputs
(
1
);
streamConfig
.
set
TypeSerializersIn
(
inputSerializer
);
streamConfig
.
set
upNetworkInputs
(
inputSerializer
);
}
public
<
K
>
void
configureForKeyedStream
(
...
...
flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/StreamConfigChainer.java
浏览文件 @
0a66d09b
...
...
@@ -144,7 +144,7 @@ public class StreamConfigChainer<OWNER> {
tailConfig
=
new
StreamConfig
(
new
Configuration
());
tailConfig
.
setStreamOperatorFactory
(
checkNotNull
(
operatorFactory
));
tailConfig
.
setOperatorID
(
checkNotNull
(
operatorID
));
tailConfig
.
set
TypeSerializersIn
(
inputSerializer
);
tailConfig
.
set
upNetworkInputs
(
inputSerializer
);
tailConfig
.
setTypeSerializerOut
(
outputSerializer
);
if
(
createKeyedStateBackend
)
{
// used to test multiple stateful operators chained in a single task.
...
...
flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/TwoInputStreamTaskTestHarness.java
浏览文件 @
0a66d09b
...
...
@@ -179,7 +179,7 @@ public class TwoInputStreamTaskTestHarness<IN1, IN2, OUT> extends StreamTaskTest
streamConfig
.
setInPhysicalEdges
(
inPhysicalEdges
);
streamConfig
.
setNumberOfNetworkInputs
(
numInputGates
);
streamConfig
.
set
TypeSerializersIn
(
inputSerializer1
,
inputSerializer2
);
streamConfig
.
set
upNetworkInputs
(
inputSerializer1
,
inputSerializer2
);
}
@Override
...
...
flink-streaming-java/src/test/java/org/apache/flink/streaming/util/OneInputStreamOperatorTestHarness.java
浏览文件 @
0a66d09b
...
...
@@ -58,7 +58,7 @@ public class OneInputStreamOperatorTestHarness<IN, OUT>
throws
Exception
{
this
(
operator
,
1
,
1
,
0
);
config
.
set
TypeSerializersIn
(
Preconditions
.
checkNotNull
(
typeSerializerIn
));
config
.
set
upNetworkInputs
(
Preconditions
.
checkNotNull
(
typeSerializerIn
));
}
public
OneInputStreamOperatorTestHarness
(
...
...
@@ -75,7 +75,7 @@ public class OneInputStreamOperatorTestHarness<IN, OUT>
parallelism
,
subtaskIndex
,
operatorID
);
config
.
set
TypeSerializersIn
(
Preconditions
.
checkNotNull
(
typeSerializerIn
));
config
.
set
upNetworkInputs
(
Preconditions
.
checkNotNull
(
typeSerializerIn
));
}
public
OneInputStreamOperatorTestHarness
(
...
...
@@ -85,7 +85,7 @@ public class OneInputStreamOperatorTestHarness<IN, OUT>
throws
Exception
{
this
(
operator
,
environment
);
config
.
set
TypeSerializersIn
(
Preconditions
.
checkNotNull
(
typeSerializerIn
));
config
.
set
upNetworkInputs
(
Preconditions
.
checkNotNull
(
typeSerializerIn
));
}
public
OneInputStreamOperatorTestHarness
(
OneInputStreamOperator
<
IN
,
OUT
>
operator
)
...
...
@@ -139,7 +139,7 @@ public class OneInputStreamOperatorTestHarness<IN, OUT>
throws
Exception
{
this
(
factory
,
environment
);
config
.
set
TypeSerializersIn
(
Preconditions
.
checkNotNull
(
typeSerializerIn
));
config
.
set
upNetworkInputs
(
Preconditions
.
checkNotNull
(
typeSerializerIn
));
}
public
OneInputStreamOperatorTestHarness
(
...
...
@@ -153,7 +153,7 @@ public class OneInputStreamOperatorTestHarness<IN, OUT>
throws
Exception
{
this
(
factory
,
1
,
1
,
0
);
config
.
set
TypeSerializersIn
(
Preconditions
.
checkNotNull
(
typeSerializerIn
));
config
.
set
upNetworkInputs
(
Preconditions
.
checkNotNull
(
typeSerializerIn
));
}
public
OneInputStreamOperatorTestHarness
(
...
...
flink-table/flink-table-runtime-blink/src/main/java/org/apache/flink/table/runtime/operators/multipleinput/MultipleInputStreamOperatorBase.java
浏览文件 @
0a66d09b
...
...
@@ -285,7 +285,7 @@ public abstract class MultipleInputStreamOperatorBase extends AbstractStreamOper
streamConfig
.
setOperatorName
(
wrapper
.
getOperatorName
());
streamConfig
.
setNumberOfNetworkInputs
(
wrapper
.
getAllInputTypes
().
size
());
streamConfig
.
setNumberOfOutputs
(
wrapper
.
getOutputEdges
().
size
());
streamConfig
.
set
TypeSerializersIn
(
streamConfig
.
set
upNetworkInputs
(
wrapper
.
getAllInputTypes
().
stream
()
.
map
(
t
->
t
.
createSerializer
(
executionConfig
))
.
toArray
(
TypeSerializer
[]::
new
));
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录