- 11 9月, 2015 1 次提交
-
-
由 Aljoscha Krettek 提交于
Before, when one source closes early it will not emit watermarks anymore. Downstream operations don't know about this and expect watermarks to keep on coming. This leads to watermarks not being forwarded anymore. Now, when a source closes it will emit a final watermark with timestamp Long.MAX_VALUE. This will have the effect of allowing the watermarks from the other operations to propagate though because the watermark is defined as the minimum over all inputs. The Long.MAX_VALUE watermark has the added benefit of notifying operations that no more elements will arrive in the future. This closes #1060
-
- 10 9月, 2015 11 次提交
-
-
由 mjsax 提交于
This closes #1114
-
由 Robert Metzger 提交于
This closes #941
-
由 Fabian Hueske 提交于
This closes #1111
-
由 Maximilian Michels 提交于
[FLINK-2645] [jobmanager] Fail job execution if final accumulators cannot be merged and forward exceptions. This closes #1112
-
由 Till Rohrmann 提交于
[FLINK-2631] [streaming] Fixes the StreamFold operator and adds OutputTypeConfigurable interface to support type injection at StreamGraph creation. Adds test for non serializable fold type. Adds test to verify proper output type forwarding for OutputTypeConfigurable implementations. Makes OutputTypeConfigurable typed, tests that TwoInputStreamOperator is output type configurable This closes #1101
-
由 Aljoscha Krettek 提交于
-
由 Stephan Ewen 提交于
-
由 Theodore Vasiloudis 提交于
This closes #1100
-
由 Stephan Ewen 提交于
This closes #1109
-
由 Stephan Ewen 提交于
-
由 Chiwan Park 提交于
This closes #1103
-
- 09 9月, 2015 11 次提交
-
-
由 Tamara Mendt 提交于
..capable of tracking min, max and estimates for count distinct and heavy hitters. The count distinct algorithms are Linear Counting and HyperLogLog, both from an imported library from clearspring. The heavy hitters algorithms are Lossy counting (Manku et.al 2002) and one based on Count Min Sketch (Cormode 2005). The heavy hitters algorithms are implemented in the statistics package in the flink-operator-stats submodule of flink-contrib. Include tests verifying if merged skecthes have the same guarantees as local sketches. Added conditional to every collect variable, to make stats more configurable. Made print operator stats more informative. Added arguments to yarn pom to enable building Implemented deep cloning for class OperatorStatistics. Added additional tests in OperatorStatsAccumulatorTest to check if the accumulator works when only one statistic is being tracked, rather than all. This closes #605.
-
由 chenliang613 提交于
This closes #1094
-
由 Vimal 提交于
This closes #1091
-
由 Aljoscha Krettek 提交于
This also bumps surefire/failsafe version to 2.8.1
-
由 Theodore Vasiloudis 提交于
This closes #889.
-
由 Robert Metzger 提交于
-
由 Stephan Ewen 提交于
[FLINK-2580] [runtime] Expose more methods form Hadoop output streams and exposes wrapped input and output streams.
-
由 Stephan Ewen 提交于
-
由 Stephan Ewen 提交于
-
由 Stephan Ewen 提交于
This closes #1093
-
由 HuangWHWHW 提交于
This closes #992.
-
- 08 9月, 2015 5 次提交
-
-
由 Till Rohrmann 提交于
-
由 r-pogalz 提交于
This closes #1052
-
由 r-pogalz 提交于
-
由 HuangWHWHW 提交于
- improve test layout This closes #1073.
-
由 Andra Lungu 提交于
[FLINK-2570] [gelly] Added a description of the I/O This closes #1054
-
- 07 9月, 2015 7 次提交
-
-
由 tammymendt 提交于
[FLINK-2567] [core] Allow quoted strings in CSV fields to contain quotation character inside of the field, as long as its escaped Ex: 'Hi my name is \'Flink\'' This closes #1059
-
由 Robert Metzger 提交于
This closes #1082
-
由 Stephan Ewen 提交于
-
由 HuangWHWHW 提交于
This closes #1096
-
由 Stephan Ewen 提交于
-
由 Stephan Ewen 提交于
The test is inconclusive when the test failure happens before the first checkpoint.
-
由 Greg Hogan 提交于
LocalExecutor and Client now pass Configuration to JobGraphGenerator. This closes #1095
-
- 05 9月, 2015 1 次提交
-
-
由 tedyu 提交于
This closes #1089.
-
- 04 9月, 2015 1 次提交
-
-
由 Nikolaas Steenbergen 提交于
[FLINK-2161] [scala shell] Modify start script to take additional argument (-a <path/to/class> or --addclasspath <path/to/class>) for external libraries This closes #805
-
- 03 9月, 2015 3 次提交
-
-
由 mjsax 提交于
This closes #1074
-
由 Maximilian Michels 提交于
- call new start() method of FlinkMiniCluster
-
由 Maximilian Michels 提交于
This closes #1085.
-