diff --git a/bin/pig b/bin/pig index cdf436a383..75901eb78e 100755 --- a/bin/pig +++ b/bin/pig @@ -52,12 +52,15 @@ case "`uname`" in CYGWIN*) cygwin=true;; esac debug=false +SPORK=false remaining=() includeHCatalog=""; addJarString=-Dpig.additional.jars\=; additionalJars=""; # filter command line parameter + + for f in "$@"; do if [[ $f == "-secretDebugCmd" || $f == "-printCmdDebug" ]]; then debug=true @@ -73,6 +76,10 @@ for f in "$@"; do else remaining[${#remaining[@]}]="$f" fi + ####for spark mode##### + if [[ ${remaining[@]} == "-x spark" ]]; then #### + SPORK=true + fi done # resolve links - $0 may be a softlink @@ -340,7 +347,8 @@ if [ "$includeHCatalog" == "true" ]; then fi # run it -if [ -n "$HADOOP_BIN" ]; then +#####modified for spark mode +if [ -n "$HADOOP_BIN" ] && [ "$SPORK" == "false" ]; then if [ "$debug" == "true" ]; then echo "Find hadoop at $HADOOP_BIN" fi @@ -356,6 +364,7 @@ if [ -n "$HADOOP_BIN" ]; then PIG_JAR=`echo $PIG_HOME/share/pig/pig-*withouthadoop.jar` fi + if [ -n "$PIG_JAR" ]; then CLASSPATH=${CLASSPATH}:$PIG_JAR else @@ -370,7 +379,6 @@ if [ -n "$HADOOP_BIN" ]; then echo "HADOOP_CLASSPATH: $HADOOP_CLASSPATH" echo "HADOOP_OPTS: $HADOOP_OPTS" echo "$HADOOP_BIN" jar "$PIG_JAR" "${remaining[@]}" - echo else exec "$HADOOP_BIN" jar "$PIG_JAR" "${remaining[@]}" fi diff --git a/ivy.xml b/ivy.xml index 63554deb85..0af2ea12e7 100644 --- a/ivy.xml +++ b/ivy.xml @@ -278,10 +278,11 @@ - + + - + diff --git a/pig-spark b/pig-spark index 0ccdcbf9de..71e6ba3719 100644 --- a/pig-spark +++ b/pig-spark @@ -18,12 +18,16 @@ export SPARK_YARN_APP_JAR=pig/pig-withouthadoop.jar export SPARK_JAR=spark.jar export SPARK_MASTER=yarn-client +export SPARK_HOME= + # jars to ship, pig-withouthadoop.jar to workaround Classloader issue export SPARK_JARS=piggybank.jar # Pig settings export PIG_CLASSPATH=${SPARK_JAR}:${SPARK_JARS}:mesos.jar:pig-withouthadoop.jar export SPARK_PIG_JAR=pig/pig-withouthadoop.jar +export PIG_CLASSPATH=${SPARK_JAR}:${SPARK_JARS}:mesos.jar:pig-withouthadoop.jar +export SPARK_PIG_JAR=pig-withouthadoop.jar # Cluster settings export SPARK_WORKER_CORES=4 diff --git a/src/org/apache/pig/backend/hadoop/executionengine/spark/SparkUtil.java b/src/org/apache/pig/backend/hadoop/executionengine/spark/SparkUtil.java index a546d639f9..57347c1007 100644 --- a/src/org/apache/pig/backend/hadoop/executionengine/spark/SparkUtil.java +++ b/src/org/apache/pig/backend/hadoop/executionengine/spark/SparkUtil.java @@ -13,8 +13,13 @@ import scala.Product2; import scala.collection.JavaConversions; import scala.collection.Seq; -import scala.reflect.ClassManifest; -import scala.reflect.ClassManifest$; + +//==============FFFFFF================= +//import scala.reflect.ClassManifest; +//import scala.reflect.ClassManifest$; +import scala.reflect.ClassTag; +import scala.reflect.ClassTag$; + import org.apache.spark.rdd.RDD; import java.io.IOException; @@ -25,18 +30,18 @@ */ public class SparkUtil { - public static ClassManifest getManifest(Class clazz) { - return ClassManifest$.MODULE$.fromClass(clazz); + public static ClassTag getManifest(Class clazz) { + return ClassTag$.MODULE$.apply(clazz); } @SuppressWarnings("unchecked") - public static ClassManifest> getTuple2Manifest() { - return (ClassManifest>)(Object)getManifest(Tuple2.class); + public static ClassTag> getTuple2Manifest() { + return (ClassTag>)(Object)getManifest(Tuple2.class); } @SuppressWarnings("unchecked") - public static ClassManifest> getProduct2Manifest() { - return (ClassManifest>)(Object)getManifest(Product2.class); + public static ClassTag> getProduct2Manifest() { + return (ClassTag>)(Object)getManifest(Product2.class); } public static JobConf newJobConf(PigContext pigContext) throws IOException { diff --git a/src/org/apache/pig/backend/hadoop/executionengine/spark/converter/DistinctConverter.java b/src/org/apache/pig/backend/hadoop/executionengine/spark/converter/DistinctConverter.java index 86ef96832a..c8ebcfd6b2 100644 --- a/src/org/apache/pig/backend/hadoop/executionengine/spark/converter/DistinctConverter.java +++ b/src/org/apache/pig/backend/hadoop/executionengine/spark/converter/DistinctConverter.java @@ -13,7 +13,10 @@ import scala.Function1; import scala.Function2; import scala.Tuple2; -import scala.reflect.ClassManifest; + +//import scala.reflect.ClassManifest; +import scala.reflect.ClassTag; + import scala.runtime.AbstractFunction1; import scala.runtime.AbstractFunction2; import org.apache.spark.rdd.PairRDDFunctions; @@ -33,7 +36,7 @@ public RDD convert(List> predecessors, PODistinct poDistinct) SparkUtil.assertPredecessorSize(predecessors, poDistinct, 1); RDD rdd = predecessors.get(0); - ClassManifest> tuple2ClassManifest = SparkUtil.getTuple2Manifest(); + ClassTag> tuple2ClassManifest = SparkUtil.getTuple2Manifest(); RDD> rddPairs = rdd.map(TO_KEY_VALUE_FUNCTION, tuple2ClassManifest); PairRDDFunctions pairRDDFunctions = diff --git a/src/org/apache/pig/backend/hadoop/executionengine/spark/converter/GlobalRearrangeConverter.java b/src/org/apache/pig/backend/hadoop/executionengine/spark/converter/GlobalRearrangeConverter.java index a386954675..e9e671c392 100644 --- a/src/org/apache/pig/backend/hadoop/executionengine/spark/converter/GlobalRearrangeConverter.java +++ b/src/org/apache/pig/backend/hadoop/executionengine/spark/converter/GlobalRearrangeConverter.java @@ -22,7 +22,9 @@ import scala.Tuple2; import scala.collection.JavaConversions; import scala.collection.Seq; -import scala.reflect.ClassManifest; + +//import scala.reflect.ClassManifest; +import scala.reflect.ClassTag; @SuppressWarnings({ "serial"}) public class GlobalRearrangeConverter implements POConverter { @@ -59,7 +61,7 @@ public RDD convert(List> predecessors, } else { //COGROUP // each pred returns (index, key, value) - ClassManifest> tuple2ClassManifest = SparkUtil.getTuple2Manifest(); + ClassTag> tuple2ClassManifest = SparkUtil.getTuple2Manifest(); List>> rddPairs = new ArrayList(); for (RDD rdd : predecessors) {