Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
doujutun3207
flink
提交
515ad3c3
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,发现更多精彩内容 >>
提交
515ad3c3
编写于
6月 20, 2014
作者:
T
Till Rohrmann
提交者:
uce
6月 25, 2014
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[FLINK-960] Fix CollectionDataSource bug
This closes #33.
上级
3d6cc5f4
变更
3
隐藏空白更改
内联
并排
Showing
3 changed file
with
59 addition
and
6 deletion
+59
-6
stratosphere-java/src/main/java/eu/stratosphere/api/java/record/operators/CollectionDataSource.java
...phere/api/java/record/operators/CollectionDataSource.java
+7
-6
stratosphere-scala/pom.xml
stratosphere-scala/pom.xml
+6
-0
stratosphere-scala/src/test/scala/eu/stratosphere/api/scala/CollectionDataSourceTest.scala
.../eu/stratosphere/api/scala/CollectionDataSourceTest.scala
+46
-0
未找到文件。
stratosphere-java/src/main/java/eu/stratosphere/api/java/record/operators/CollectionDataSource.java
浏览文件 @
515ad3c3
...
...
@@ -136,13 +136,14 @@ public class CollectionDataSource extends GenericDataSourceBase<Record, GenericI
checkFormat
((
Collection
<
Object
>)
data
[
0
]);
f
.
setData
((
Collection
<
Object
>)
data
[
0
]);
}
Collection
<
Object
>
tmp
=
new
ArrayList
<
Object
>();
for
(
Object
o
:
data
)
{
tmp
.
add
(
o
);
else
{
Collection
<
Object
>
tmp
=
new
ArrayList
<
Object
>();
for
(
Object
o
:
data
)
{
tmp
.
add
(
o
);
}
checkFormat
(
tmp
);
f
.
setData
(
tmp
);
}
checkFormat
(
tmp
);
f
.
setData
(
tmp
);
}
// --------------------------------------------------------------------------------------------
...
...
stratosphere-scala/pom.xml
浏览文件 @
515ad3c3
...
...
@@ -56,6 +56,12 @@
<artifactId>
asm
</artifactId>
<version>
4.0
</version>
</dependency>
<dependency>
<groupId>
org.scalatest
</groupId>
<artifactId>
scalatest_2.10
</artifactId>
<version>
2.2.0
</version>
</dependency>
</dependencies>
<build>
...
...
stratosphere-scala/src/test/scala/eu/stratosphere/api/scala/CollectionDataSourceTest.scala
0 → 100644
浏览文件 @
515ad3c3
/*
* Copyright (C) 2010-2013 by the Stratosphere project (http://stratosphere.eu)
*
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on
* an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the
* specific language governing permissions and limitations under the License.
*/
package
eu.stratosphere.api.scala
import
eu.stratosphere.api.java.record.operators.
{
CollectionDataSource
=>
JCollectionDataSource
}
import
eu.stratosphere.types.
{
DoubleValue
,
Record
}
import
org.scalatest.junit.AssertionsForJUnit
import
org.junit.Assert._
import
org.junit.Test
class
CollectionDataSourceTest
extends
AssertionsForJUnit
{
@Test
def
testScalaCollectionInput
()
{
val
expected
=
List
(
1.0
,
2.0
,
3.0
)
val
datasource
=
CollectionDataSource
(
expected
)
val
javaCDS
=
datasource
.
contract
.
asInstanceOf
[
JCollectionDataSource
]
val
inputFormat
=
javaCDS
.
getFormatWrapper
.
getUserCodeObject
()
val
splits
=
inputFormat
.
createInputSplits
(
1
)
inputFormat
.
open
(
splits
(
0
))
val
record
=
new
Record
()
var
result
=
List
[
Double
]()
while
(!
inputFormat
.
reachedEnd
()){
inputFormat
.
nextRecord
(
record
)
assertTrue
(
record
.
getNumFields
==
1
)
val
value
=
record
.
getField
[
DoubleValue
](
0
,
classOf
[
DoubleValue
])
result
=
value
.
getValue
::
result
}
assertEquals
(
expected
,
result
.
reverse
)
}
}
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录