This is an automated email from the ASF dual-hosted git repository. rpardomeza pushed a commit to branch blossom-teradata in repository https://gitbox.apache.org/repos/asf/incubator-wayang.git
commit cb47c12191b825872347d8108c9e3e3efc846c40 Author: Rodrigo Pardo Meza <[email protected]> AuthorDate: Tue Jun 14 23:17:23 2022 +0200 [TERADATA] Base Teradata source access test to be deleted --- .../org/apache/wayang/teradata/TestTeradata.java | 120 +++++++++++++++++++++ 1 file changed, 120 insertions(+) diff --git a/wayang-plugins/wayang-teradata/src/main/java/org/apache/wayang/teradata/TestTeradata.java b/wayang-plugins/wayang-teradata/src/main/java/org/apache/wayang/teradata/TestTeradata.java new file mode 100644 index 00000000..0d4ecee4 --- /dev/null +++ b/wayang-plugins/wayang-teradata/src/main/java/org/apache/wayang/teradata/TestTeradata.java @@ -0,0 +1,120 @@ +package org.apache.wayang.teradata; + +import org.apache.wayang.basic.data.Record; +import org.apache.wayang.basic.data.Tuple2; +import org.apache.wayang.basic.operators.*; +import org.apache.wayang.basic.types.RecordType; +import org.apache.wayang.core.api.WayangContext; +import org.apache.wayang.core.function.FlatMapDescriptor; +import org.apache.wayang.core.function.ReduceDescriptor; +import org.apache.wayang.core.function.TransformationDescriptor; +import org.apache.wayang.core.optimizer.ProbabilisticDoubleInterval; +import org.apache.wayang.core.plan.wayangplan.WayangPlan; +import org.apache.wayang.core.types.DataSetType; +import org.apache.wayang.core.types.DataUnitType; +import org.apache.wayang.core.util.ReflectionUtils; +import org.apache.wayang.java.Java; +import org.apache.wayang.java.platform.JavaPlatform; +import org.apache.wayang.teradata.operator.TeradataTableOperator; + +import java.util.Arrays; +import java.util.LinkedList; +import java.util.List; + +public class TestTeradata { + + public static void main(String[] args) { + + List<Tuple2<String, Integer>> collector = new LinkedList<>(); + TeradataTableOperator table = new TeradataTableOperator("test.employees"); + table.setName("Table"); + List<Record> collector2 = new LinkedList<>(); + + LocalCallbackSink<Record> sink = LocalCallbackSink.createCollectingSink(collector2, DataSetType.createDefault(Record.class)); + sink.setName("Collect result"); + + // Build Rheem plan by connecting operators + table.connectTo(0, sink, 0); + + WayangContext wayangContext = new WayangContext(); + wayangContext.register(Java.basicPlugin()); + wayangContext.register(Teradata.plugin()); + + wayangContext.execute( + new WayangPlan(sink), + ReflectionUtils.getDeclaringJar(TestTeradata.class), + ReflectionUtils.getDeclaringJar(JavaPlatform.class)); + + collector2.forEach(t -> System.out.println(t.toString())); + + } + + public static void main_wc(){ + List<Tuple2<String, Integer>> collector = new LinkedList<>(); + TextFileSource textFileSource = new TextFileSource("file:///Users/rodrigopardomeza/files/demo.txt"); + textFileSource.setName("Load file"); + + // for each line (input) output an iterator of the words + FlatMapOperator<String, String> flatMapOperator = new FlatMapOperator<>( + new FlatMapDescriptor<>(line -> Arrays.asList(line.split("\\W+")), + String.class, + String.class, + new ProbabilisticDoubleInterval(100, 10000, 0.8) + ) + ); + flatMapOperator.setName("Split words"); + + FilterOperator<String> filterOperator = new FilterOperator<>(str -> !str.isEmpty(), String.class); + filterOperator.setName("Filter empty words"); + + + // for each word transform it to lowercase and output a key-value pair (word, 1) + MapOperator<String, Tuple2<String, Integer>> mapOperator = new MapOperator<>( + new TransformationDescriptor<>(word -> new Tuple2<>(word.toLowerCase(), 1), + DataUnitType.createBasic(String.class), + DataUnitType.createBasicUnchecked(Tuple2.class) + ), DataSetType.createDefault(String.class), + DataSetType.createDefaultUnchecked(Tuple2.class) + ); + mapOperator.setName("To lower case, add counter"); + + + // groupby the key (word) and add up the values (frequency) + ReduceByOperator<Tuple2<String, Integer>, String> reduceByOperator = new ReduceByOperator<>( + new TransformationDescriptor<>(pair -> pair.field0, + DataUnitType.createBasicUnchecked(Tuple2.class), + DataUnitType.createBasic(String.class)), new ReduceDescriptor<>( + ((a, b) -> { + a.field1 += b.field1; + return a; + }), DataUnitType.createGroupedUnchecked(Tuple2.class), + DataUnitType.createBasicUnchecked(Tuple2.class) + ), DataSetType.createDefaultUnchecked(Tuple2.class) + ); + reduceByOperator.setName("Add counters"); + + + // write results to a sink + LocalCallbackSink<Tuple2<String, Integer>> sink = LocalCallbackSink.createCollectingSink( + collector, + DataSetType.createDefaultUnchecked(Tuple2.class) + ); + sink.setName("Collect result"); + + // Build Rheem plan by connecting operators + textFileSource.connectTo(0, flatMapOperator, 0); + flatMapOperator.connectTo(0, filterOperator, 0); + filterOperator.connectTo(0, mapOperator, 0); + mapOperator.connectTo(0, reduceByOperator, 0); + reduceByOperator.connectTo(0, sink, 0); + + WayangContext wayangContext = new WayangContext(); + wayangContext.register(Java.basicPlugin()); + + wayangContext.execute(new WayangPlan(sink), ReflectionUtils.getDeclaringJar(TestTeradata.class), ReflectionUtils.getDeclaringJar(JavaPlatform.class)); + + collector.sort((t1, t2) -> Integer.compare(t2.field1, t1.field1)); + System.out.printf("Found %d words:\n", collector.size()); + collector.forEach(wc -> System.out.printf("%dx %s\n", wc.field1, wc.field0)); + } +}
