This is an automated email from the ASF dual-hosted git repository.
zkaoudi pushed a change to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-wayang.git
from 9882b860 Merge pull request #517 from siddhant-grfrs/patch-1
add a6e84d03 Fix distributed usage of wayang-flink, enabling real use on
Flink clusters
add f7cf47f2 Remove unused serialization tools
add de5e6899 Remove logging and old comments
add a3e7fe1f Remove setting parallelism on unions as its not allowed
add a5e9633d Ignore non-working tensorflow integration test
add 96558df9 Merge pull request #520 from
juripetersen/fix-flink-distributed
No new revisions were added by this update.
Summary of changes:
pom.xml | 2 +-
.../wayang/core/plan/wayangplan/Operator.java | 23 ++++++-
wayang-platforms/wayang-flink/pom.xml | 20 ++++++
.../wayang/flink/compiler/FunctionCompiler.java | 1 -
.../wayang/flink/compiler/KeySelectorFunction.java | 5 +-
.../wayang/flink/execution/FlinkExecutor.java | 5 +-
.../operators/CollectionSplittableIterator.java | 74 ++++++++++++++++++++++
.../flink/operators/FlinkCartesianOperator.java | 4 +-
.../flink/operators/FlinkCoGroupOperator.java | 6 +-
.../flink/operators/FlinkCollectionSink.java | 18 +++++-
.../flink/operators/FlinkCollectionSource.java | 52 +++++++++++++--
.../flink/operators/FlinkDistinctOperator.java | 3 +-
.../flink/operators/FlinkFilterOperator.java | 3 +-
.../flink/operators/FlinkFlatMapOperator.java | 5 +-
.../flink/operators/FlinkGlobalReduceOperator.java | 3 +-
.../flink/operators/FlinkGroupByOperator.java | 5 +-
.../wayang/flink/operators/FlinkJoinOperator.java | 11 ++--
.../flink/operators/FlinkLocalCallbackSink.java | 5 +-
.../wayang/flink/operators/FlinkMapOperator.java | 10 +--
.../operators/FlinkMapPartitionsOperator.java | 7 +-
.../flink/operators/FlinkObjectFileSink.java | 6 +-
.../flink/operators/FlinkObjectFileSource.java | 7 +-
.../flink/operators/FlinkReduceByOperator.java | 2 +-
.../operators/FlinkRepeatExpandedOperator.java | 4 +-
.../flink/operators/FlinkRepeatOperator.java | 6 +-
.../wayang/flink/operators/FlinkSortOperator.java | 3 +-
.../wayang/flink/operators/FlinkTextFileSink.java | 6 +-
.../flink/operators/FlinkTextFileSource.java | 2 +-
.../wayang/flink/operators/FlinkTsvFileSink.java | 9 ++-
.../flink/operators/ScalaTupleSerializer.java | 61 ++++++++++++++++++
.../wayang/flink/platform/FlinkPlatform.java | 16 +++--
.../resources/wayang-flink-defaults.properties | 2 +-
.../org/apache/wayang/tests/TensorflowIrisIT.java | 2 +-
33 files changed, 327 insertions(+), 61 deletions(-)
create mode 100644
wayang-platforms/wayang-flink/src/main/java/org/apache/wayang/flink/operators/CollectionSplittableIterator.java
create mode 100644
wayang-platforms/wayang-flink/src/main/java/org/apache/wayang/flink/operators/ScalaTupleSerializer.java