Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
doujutun3207
flink
提交
937412cb
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,发现更多精彩内容 >>
提交
937412cb
编写于
7月 14, 2014
作者:
M
Márton Balassi
提交者:
Stephan Ewen
8月 18, 2014
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[streaming] Package naming refactor
上级
8c219335
变更
15
隐藏空白更改
内联
并排
Showing
15 changed file
with
36 addition
and
20 deletion
+36
-20
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/api/JobGraphBuilder.java
...n/java/eu/stratosphere/streaming/api/JobGraphBuilder.java
+5
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/api/StreamSink.java
...c/main/java/eu/stratosphere/streaming/api/StreamSink.java
+3
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/api/StreamSource.java
...main/java/eu/stratosphere/streaming/api/StreamSource.java
+4
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/api/StreamTask.java
...c/main/java/eu/stratosphere/streaming/api/StreamTask.java
+3
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/api/invokable/DefaultSourceInvokable.java
...phere/streaming/api/invokable/DefaultSourceInvokable.java
+1
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/api/invokable/DefaultTaskInvokable.java
...osphere/streaming/api/invokable/DefaultTaskInvokable.java
+1
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/api/invokable/UserSinkInvokable.java
...ratosphere/streaming/api/invokable/UserSinkInvokable.java
+1
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/api/invokable/UserSourceInvokable.java
...tosphere/streaming/api/invokable/UserSourceInvokable.java
+1
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/api/invokable/UserTaskInvokable.java
...ratosphere/streaming/api/invokable/UserTaskInvokable.java
+1
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/test/MyStream.java
...rc/main/java/eu/stratosphere/streaming/test/MyStream.java
+7
-6
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/test/QuerySourceInvokable.java
.../eu/stratosphere/streaming/test/QuerySourceInvokable.java
+2
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/test/RandIS.java
.../src/main/java/eu/stratosphere/streaming/test/RandIS.java
+1
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/test/TestSinkInvokable.java
...ava/eu/stratosphere/streaming/test/TestSinkInvokable.java
+2
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/test/TestSourceInvokable.java
...a/eu/stratosphere/streaming/test/TestSourceInvokable.java
+2
-1
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/test/TestTaskInvokable.java
...ava/eu/stratosphere/streaming/test/TestTaskInvokable.java
+2
-1
未找到文件。
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/JobGraphBuilder.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
api/
JobGraphBuilder.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.api
;
import
java.util.HashMap
;
import
java.util.Map
;
...
...
@@ -14,6 +14,9 @@ import eu.stratosphere.nephele.jobgraph.JobOutputVertex;
import
eu.stratosphere.nephele.jobgraph.JobTaskVertex
;
import
eu.stratosphere.nephele.template.AbstractInputTask
;
import
eu.stratosphere.pact.runtime.task.util.TaskConfig
;
import
eu.stratosphere.streaming.api.invokable.UserSinkInvokable
;
import
eu.stratosphere.streaming.api.invokable.UserSourceInvokable
;
import
eu.stratosphere.streaming.api.invokable.UserTaskInvokable
;
import
eu.stratosphere.streaming.partitioner.DefaultPartitioner
;
import
eu.stratosphere.types.Record
;
...
...
@@ -89,6 +92,7 @@ public class JobGraphBuilder {
components
.
put
(
sinkName
,
sink
);
}
//TODO refactor connects
public
void
connect
(
String
upStreamComponentName
,
String
downStreamComponentName
,
ChannelType
channelType
)
{
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/StreamSink.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
api/
StreamSink.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.api
;
import
eu.stratosphere.nephele.io.RecordReader
;
import
eu.stratosphere.nephele.template.AbstractOutputTask
;
import
eu.stratosphere.streaming.api.invokable.UserSinkInvokable
;
import
eu.stratosphere.streaming.test.TestSinkInvokable
;
import
eu.stratosphere.types.Record
;
public
class
StreamSink
extends
AbstractOutputTask
{
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/StreamSource.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
api/
StreamSource.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.api
;
import
eu.stratosphere.nephele.io.ChannelSelector
;
import
eu.stratosphere.nephele.io.RecordWriter
;
import
eu.stratosphere.nephele.template.AbstractInputTask
;
import
eu.stratosphere.streaming.api.invokable.UserSourceInvokable
;
import
eu.stratosphere.streaming.partitioner.DefaultPartitioner
;
import
eu.stratosphere.streaming.test.RandIS
;
import
eu.stratosphere.streaming.test.TestSourceInvokable
;
import
eu.stratosphere.types.Record
;
public
class
StreamSource
extends
AbstractInputTask
<
RandIS
>
{
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/StreamTask.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
api/
StreamTask.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.api
;
import
java.util.LinkedList
;
import
java.util.List
;
...
...
@@ -7,7 +7,9 @@ import eu.stratosphere.nephele.io.ChannelSelector;
import
eu.stratosphere.nephele.io.RecordReader
;
import
eu.stratosphere.nephele.io.RecordWriter
;
import
eu.stratosphere.nephele.template.AbstractTask
;
import
eu.stratosphere.streaming.api.invokable.UserTaskInvokable
;
import
eu.stratosphere.streaming.partitioner.DefaultPartitioner
;
import
eu.stratosphere.streaming.test.TestTaskInvokable
;
import
eu.stratosphere.types.Record
;
public
class
StreamTask
extends
AbstractTask
{
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/DefaultSourceInvokable.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
api/invokable/
DefaultSourceInvokable.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.api.invokable
;
import
eu.stratosphere.nephele.io.RecordWriter
;
import
eu.stratosphere.types.Record
;
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/DefaultTaskInvokable.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
api/invokable/
DefaultTaskInvokable.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.api.invokable
;
import
eu.stratosphere.nephele.io.RecordWriter
;
import
eu.stratosphere.types.Record
;
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/UserSinkInvokable.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
api/invokable/
UserSinkInvokable.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.api.invokable
;
import
eu.stratosphere.nephele.io.RecordReader
;
import
eu.stratosphere.types.Record
;
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/UserSourceInvokable.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
api/invokable/
UserSourceInvokable.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.api.invokable
;
import
eu.stratosphere.nephele.io.RecordWriter
;
import
eu.stratosphere.types.Record
;
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/UserTaskInvokable.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
api/invokable/
UserTaskInvokable.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.api.invokable
;
import
eu.stratosphere.nephele.io.RecordWriter
;
import
eu.stratosphere.types.Record
;
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/MyStream.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
test/
MyStream.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.test
;
import
eu.stratosphere.nephele.io.channels.ChannelType
;
import
eu.stratosphere.nephele.jobgraph.JobGraph
;
import
eu.stratosphere.streaming.JobGraphBuilder.Partitioning
;
import
eu.stratosphere.streaming.api.JobGraphBuilder
;
import
eu.stratosphere.streaming.api.JobGraphBuilder.Partitioning
;
import
eu.stratosphere.test.util.TestBase2
;
public
class
MyStream
extends
TestBase2
{
...
...
@@ -15,10 +16,10 @@ public class MyStream extends TestBase2 {
graphBuilder
.
setTask
(
"cellTask"
,
TestTaskInvokable
.
class
,
Partitioning
.
BROADCAST
,
2
);
graphBuilder
.
setSink
(
"sink"
,
TestSinkInvokable
.
class
);
graphBuilder
.
connect
(
"infoSource"
,
"cellTask"
,
Partitioning
.
BROADCAST
,
ChannelType
.
INMEMORY
);
graphBuilder
.
connect
(
"querySource"
,
"cellTask"
,
Partitioning
.
BROADCAST
,
ChannelType
.
INMEMORY
);
graphBuilder
.
connect
(
"cellTask"
,
"sink"
,
Partitioning
.
BROADCAST
,
ChannelType
.
INMEMORY
);
graphBuilder
.
connect
(
"infoSource"
,
"cellTask"
,
Partitioning
.
BROADCAST
,
ChannelType
.
INMEMORY
);
graphBuilder
.
connect
(
"querySource"
,
"cellTask"
,
Partitioning
.
BROADCAST
,
ChannelType
.
INMEMORY
);
graphBuilder
.
connect
(
"cellTask"
,
"sink"
,
Partitioning
.
BROADCAST
,
ChannelType
.
INMEMORY
);
return
graphBuilder
.
getJobGraph
();
}
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/QuerySourceInvokable.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
test/
QuerySourceInvokable.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.test
;
import
eu.stratosphere.nephele.io.RecordWriter
;
import
eu.stratosphere.streaming.api.invokable.UserSourceInvokable
;
import
eu.stratosphere.types.IntValue
;
import
eu.stratosphere.types.LongValue
;
import
eu.stratosphere.types.Record
;
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/RandIS.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
test/
RandIS.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.test
;
import
java.io.DataInput
;
import
java.io.DataOutput
;
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/TestSinkInvokable.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
test/
TestSinkInvokable.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.test
;
import
eu.stratosphere.nephele.io.RecordReader
;
import
eu.stratosphere.streaming.api.invokable.UserSinkInvokable
;
import
eu.stratosphere.types.Record
;
import
eu.stratosphere.types.StringValue
;
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/TestSourceInvokable.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
test/
TestSourceInvokable.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.test
;
import
eu.stratosphere.nephele.io.RecordWriter
;
import
eu.stratosphere.streaming.api.invokable.UserSourceInvokable
;
import
eu.stratosphere.types.IntValue
;
import
eu.stratosphere.types.LongValue
;
import
eu.stratosphere.types.Record
;
...
...
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/TestTaskInvokable.java
→
flink-addons/flink-streaming/src/main/java/eu/stratosphere/streaming/
test/
TestTaskInvokable.java
浏览文件 @
937412cb
package
eu.stratosphere.streaming
;
package
eu.stratosphere.streaming
.test
;
import
eu.stratosphere.nephele.io.RecordWriter
;
import
eu.stratosphere.streaming.api.invokable.UserTaskInvokable
;
import
eu.stratosphere.streaming.cellinfo.WorkerEngineExact
;
import
eu.stratosphere.types.IntValue
;
import
eu.stratosphere.types.LongValue
;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录