Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
doujutun3207
flink
提交
7582390c
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,发现更多精彩内容 >>
提交
7582390c
编写于
3月 07, 2015
作者:
M
mbalassi
提交者:
Gábor Hermann
3月 08, 2015
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[FLINK-1429] [streaming] Scala programming guide update: intro & operators, minor fixes
This closes #463
上级
bd1b916f
变更
3
展开全部
显示空白变更内容
内联
并排
Showing
3 changed file
with
473 addition
and
146 deletion
+473
-146
docs/programming_guide.md
docs/programming_guide.md
+1
-1
docs/streaming_guide.md
docs/streaming_guide.md
+470
-126
flink-staging/flink-streaming/flink-streaming-scala/src/main/scala/org/apache/flink/streaming/api/scala/DataStream.scala
...ala/org/apache/flink/streaming/api/scala/DataStream.scala
+2
-19
未找到文件。
docs/programming_guide.md
浏览文件 @
7582390c
docs/streaming_guide.md
浏览文件 @
7582390c
此差异已折叠。
点击以展开。
flink-staging/flink-streaming/flink-streaming-scala/src/main/scala/org/apache/flink/streaming/api/scala/DataStream.scala
浏览文件 @
7582390c
...
...
@@ -221,8 +221,8 @@ class DataStream[T](javaStream: JavaStream[T]) {
*
*
*/
def
iterate
[
R
](
maxWaitTimeMillis
:
Long
=
0
)
(
stepFunction
:
DataStream
[
T
]
=>
(
DataStream
[
T
],
DataStream
[
R
]))
:
DataStream
[
R
]
=
{
def
iterate
[
R
](
maxWaitTimeMillis
:
Long
=
0
)
(
stepFunction
:
DataStream
[
T
]
=>
(
DataStream
[
T
],
DataStream
[
R
]))
:
DataStream
[
R
]
=
{
val
iterativeStream
=
javaStream
.
iterate
(
maxWaitTimeMillis
)
val
(
feedback
,
output
)
=
stepFunction
(
new
DataStream
[
T
](
iterativeStream
))
...
...
@@ -495,23 +495,6 @@ class DataStream[T](javaStream: JavaStream[T]) {
*/
def
split
(
selector
:
OutputSelector
[
T
])
:
SplitDataStream
[
T
]
=
javaStream
.
split
(
selector
)
// /**
// * Creates a new SplitDataStream that contains only the elements satisfying the
// * given output selector predicate.
// */
// def split(fun: T => String): SplitDataStream[T] = {
// if (fun == null) {
// throw new NullPointerException("OutputSelector must not be null.")
// }
// val selector = new OutputSelector[T] {
// val cleanFun = clean(fun)
// def select(in: T): java.lang.Iterable[String] = {
// List(cleanFun(in))
// }
// }
// split(selector)
// }
/**
* Creates a new SplitDataStream that contains only the elements satisfying the
* given output selector predicate.
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录