You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@beam.apache.org by ke...@apache.org on 2016/11/12 02:28:40 UTC

[22/39] incubator-beam git commit: Merge branch 'master' into apex-runner

Merge branch 'master' into apex-runner


Project: http://git-wip-us.apache.org/repos/asf/incubator-beam/repo
Commit: http://git-wip-us.apache.org/repos/asf/incubator-beam/commit/6fc47ed1
Tree: http://git-wip-us.apache.org/repos/asf/incubator-beam/tree/6fc47ed1
Diff: http://git-wip-us.apache.org/repos/asf/incubator-beam/diff/6fc47ed1

Branch: refs/heads/master
Commit: 6fc47ed101f02aacacecd5c62c4a026eaba3e12c
Parents: fa3a6aa 215980a
Author: Thomas Weise <th...@apache.org>
Authored: Thu Oct 27 19:25:06 2016 -0700
Committer: Thomas Weise <th...@apache.org>
Committed: Thu Oct 27 19:25:06 2016 -0700

----------------------------------------------------------------------
 .travis.yml                                     |  22 +-
 .../beam/examples/DebuggingWordCount.java       |   2 +-
 .../apache/beam/examples/MinimalWordCount.java  |   2 +-
 .../org/apache/beam/examples/WordCount.java     |   2 +-
 .../apache/beam/examples/complete/TfIdf.java    |   2 +-
 .../examples/complete/TopWikipediaSessions.java |   2 +-
 .../examples/cookbook/BigQueryTornadoes.java    |   2 +-
 .../cookbook/CombinePerKeyExamples.java         |   2 +-
 .../beam/examples/cookbook/DeDupExample.java    |   2 +-
 .../beam/examples/cookbook/FilterExamples.java  |   2 +-
 .../beam/examples/cookbook/JoinExamples.java    |   2 +-
 .../examples/cookbook/MaxPerKeyExamples.java    |   2 +-
 .../org/apache/beam/examples/WordCountTest.java |   2 +-
 .../examples/complete/AutoCompleteTest.java     |   6 +-
 .../beam/examples/complete/TfIdfTest.java       |   2 +-
 .../complete/TopWikipediaSessionsTest.java      |   2 +-
 .../examples/cookbook/DeDupExampleTest.java     |   4 +-
 .../examples/cookbook/JoinExamplesTest.java     |   2 +-
 .../examples/cookbook/TriggerExampleTest.java   |   2 +-
 .../beam/examples/MinimalWordCountJava8.java    |   2 +-
 .../beam/examples/complete/game/GameStats.java  |  52 +-
 .../examples/complete/game/HourlyTeamScore.java |  30 +-
 .../examples/complete/game/LeaderBoard.java     |  43 +-
 .../beam/examples/complete/game/UserScore.java  |  16 +-
 .../complete/game/utils/WriteToBigQuery.java    |  49 +-
 .../game/utils/WriteWindowedToBigQuery.java     |  14 +-
 .../examples/complete/game/GameStatsTest.java   |   2 +-
 .../complete/game/HourlyTeamScoreTest.java      |   2 +-
 .../examples/complete/game/LeaderBoardTest.java |  10 +-
 .../examples/complete/game/UserScoreTest.java   |   6 +-
 pom.xml                                         |  47 +-
 .../runners/core/GroupAlsoByWindowsDoFn.java    |  19 -
 .../runners/direct/CloningBundleFactory.java    |  98 ++++
 .../beam/runners/direct/DirectRunner.java       |   7 +-
 .../runners/direct/DoFnLifecycleManager.java    |  56 +-
 .../GroupAlsoByWindowEvaluatorFactory.java      |   8 +-
 .../direct/ImmutableListBundleFactory.java      |   4 +-
 .../beam/runners/direct/WatermarkManager.java   |  17 +-
 .../direct/WriteWithShardingFactory.java        |   6 +-
 .../direct/CloningBundleFactoryTest.java        | 177 +++++++
 .../ConsumerTrackingPipelineVisitorTest.java    |  32 +-
 .../beam/runners/direct/DirectRunnerTest.java   |  40 +-
 .../direct/DoFnLifecycleManagerTest.java        |  74 ++-
 .../EncodabilityEnforcementFactoryTest.java     |   6 +-
 .../ImmutabilityCheckingBundleFactoryTest.java  |   8 +-
 .../ImmutabilityEnforcementFactoryTest.java     |   8 +-
 .../direct/KeyedPValueTrackingVisitorTest.java  |   8 +-
 .../beam/runners/direct/ParDoEvaluatorTest.java |   8 +-
 .../direct/ParDoMultiEvaluatorFactoryTest.java  |  87 ++--
 .../direct/ParDoSingleEvaluatorFactoryTest.java |  94 ++--
 .../runners/direct/WatermarkManagerTest.java    |   8 +-
 .../dataflow/BlockingDataflowRunner.java        |  13 +-
 .../runners/dataflow/DataflowPipelineJob.java   |  17 +-
 .../dataflow/DataflowPipelineTranslator.java    |   4 +
 .../beam/runners/dataflow/DataflowRunner.java   |   4 +-
 .../options/DataflowPipelineOptions.java        |  12 +
 .../beam/runners/dataflow/util/DoFnInfo.java    |  43 +-
 runners/spark/pom.xml                           |  66 ++-
 .../runners/spark/SparkPipelineOptions.java     |  11 +
 .../apache/beam/runners/spark/SparkRunner.java  |  19 -
 .../metrics/AggregatorMetricSource.java         |   9 +-
 .../metrics/WithNamedAggregatorsSupport.java    |   7 +-
 .../coders/BeamSparkRunnerRegistrator.java      |  46 ++
 .../runners/spark/io/EmptyCheckpointMark.java   |  52 ++
 .../apache/beam/runners/spark/io/KafkaIO.java   | 131 -----
 .../beam/runners/spark/io/MicrobatchSource.java | 262 ++++++++++
 .../beam/runners/spark/io/SourceDStream.java    | 156 ++++++
 .../apache/beam/runners/spark/io/SourceRDD.java |  75 ++-
 .../runners/spark/io/SparkUnboundedSource.java  | 167 ++++++
 .../spark/stateful/StateSpecFunctions.java      | 167 ++++++
 .../runners/spark/stateful/package-info.java    |  22 +
 .../spark/translation/EvaluationContext.java    |   6 +-
 .../translation/GroupCombineFunctions.java      |  66 +--
 .../spark/translation/SparkContextFactory.java  |   5 +-
 .../spark/translation/SparkRuntimeContext.java  |   2 +-
 .../spark/translation/TransformTranslator.java  |  43 +-
 .../spark/translation/TranslationUtils.java     |  17 +
 .../SparkRunnerStreamingContextFactory.java     |   7 +
 .../streaming/StreamingEvaluationContext.java   |   6 +-
 .../streaming/StreamingTransformTranslator.java | 100 +---
 .../spark/util/SparkSideInputReader.java        |   2 +-
 .../runners/spark/ClearAggregatorsRule.java     |  33 --
 .../runners/spark/InMemoryMetricsSinkRule.java  |  32 --
 .../metrics/sink/ClearAggregatorsRule.java      |  33 ++
 .../metrics/sink/InMemoryMetrics.java           |  15 +-
 .../metrics/sink/InMemoryMetricsSinkRule.java   |  31 ++
 .../metrics/sink/NamedAggregatorsTest.java      |  92 ++++
 .../streaming/KafkaStreamingTest.java           | 150 ++++--
 .../ResumeFromCheckpointStreamingTest.java      |  34 +-
 .../streaming/utils/EmbeddedKafkaCluster.java   |  51 +-
 .../utils/KafkaWriteOnBatchCompleted.java       | 105 ++++
 .../org/apache/beam/sdk/PipelineResult.java     |  10 +-
 .../java/org/apache/beam/sdk/io/AvroIO.java     | 143 +++++-
 .../beam/sdk/options/ValueProviderUtils.java    |  60 +++
 .../sdk/transforms/DelegatingAggregator.java    |   7 +-
 .../sdk/transforms/display/DisplayData.java     |   6 +-
 .../beam/sdk/transforms/windowing/AfterAll.java |  49 --
 .../windowing/AfterDelayFromFirstElement.java   |  99 ----
 .../sdk/transforms/windowing/AfterEach.java     |  61 ---
 .../sdk/transforms/windowing/AfterFirst.java    |  50 --
 .../sdk/transforms/windowing/AfterPane.java     |  52 --
 .../windowing/AfterProcessingTime.java          |   7 -
 .../AfterSynchronizedProcessingTime.java        |   7 -
 .../transforms/windowing/AfterWatermark.java    | 158 ------
 .../transforms/windowing/DefaultTrigger.java    |  35 --
 .../beam/sdk/transforms/windowing/Never.java    |  17 -
 .../transforms/windowing/OrFinallyTrigger.java  |  46 --
 .../sdk/transforms/windowing/Repeatedly.java    |  30 --
 .../beam/sdk/transforms/windowing/Trigger.java  | 412 ++-------------
 .../apache/beam/sdk/util/ExecutableTrigger.java |  40 +-
 .../apache/beam/sdk/util/FinishedTriggers.java  |  44 --
 .../beam/sdk/util/FinishedTriggersBitSet.java   |  67 ---
 .../beam/sdk/util/FinishedTriggersSet.java      |  72 ---
 .../apache/beam/sdk/util/ReshuffleTrigger.java  |  14 -
 .../beam/sdk/util/TriggerContextFactory.java    | 507 -------------------
 .../apache/beam/sdk/util/ValueWithRecordId.java |   8 +-
 .../java/org/apache/beam/sdk/io/AvroIOTest.java |  29 ++
 .../beam/sdk/options/ValueProviderTest.java     |  17 +-
 .../sdk/options/ValueProviderUtilsTest.java     |  78 +++
 .../sdk/transforms/windowing/AfterAllTest.java  |  98 ----
 .../sdk/transforms/windowing/AfterEachTest.java |  64 ---
 .../transforms/windowing/AfterFirstTest.java    | 120 -----
 .../sdk/transforms/windowing/AfterPaneTest.java |  77 ---
 .../windowing/AfterProcessingTimeTest.java      |  94 ----
 .../AfterSynchronizedProcessingTimeTest.java    |  75 ---
 .../windowing/AfterWatermarkTest.java           | 308 -----------
 .../windowing/DefaultTriggerTest.java           | 130 -----
 .../sdk/transforms/windowing/NeverTest.java     |  34 +-
 .../windowing/OrFinallyTriggerTest.java         | 136 -----
 .../transforms/windowing/RepeatedlyTest.java    | 161 +-----
 .../sdk/transforms/windowing/StubTrigger.java   |  17 -
 .../sdk/transforms/windowing/TriggerTest.java   |  28 -
 .../beam/sdk/util/ExecutableTriggerTest.java    |  18 -
 .../sdk/util/FinishedTriggersBitSetTest.java    |  55 --
 .../sdk/util/FinishedTriggersProperties.java    | 110 ----
 .../beam/sdk/util/FinishedTriggersSetTest.java  |  60 ---
 .../beam/sdk/util/ReshuffleTriggerTest.java     |  23 -
 .../org/apache/beam/sdk/util/TriggerTester.java | 410 ---------------
 138 files changed, 2738 insertions(+), 4620 deletions(-)
----------------------------------------------------------------------