From f8ca4478923f56cf0a91394ea5cba98ba8da5241 Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Tue, 1 Dec 2015 11:34:49 -0800 Subject: [PATCH 1/8] Add predicate pushdown using filter2 api --- .../org/apache/parquet/pig/ParquetLoader.java | 148 +++++++++++++++++- .../apache/parquet/pig/TestParquetLoader.java | 53 +++++-- pom.xml | 4 +- 3 files changed, 190 insertions(+), 15 deletions(-) diff --git a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java index 41ce738b33..69227afd25 100644 --- a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java +++ b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java @@ -30,9 +30,14 @@ import static org.apache.parquet.pig.TupleReadSupport.PARQUET_COLUMN_INDEX_ACCESS; import static org.apache.parquet.pig.TupleReadSupport.getPigSchemaFromMultipleFiles; +import static org.apache.parquet.filter2.predicate.FilterApi.*; + import java.io.IOException; +import java.io.Serializable; import java.lang.ref.Reference; import java.lang.ref.SoftReference; +import java.util.ArrayList; +import java.util.Arrays; import java.util.List; import java.util.Map; import java.util.WeakHashMap; @@ -42,9 +47,15 @@ import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.RecordReader; import org.apache.hadoop.mapreduce.TaskAttemptContext; +import org.apache.parquet.filter2.predicate.FilterPredicate; +import org.apache.parquet.filter2.predicate.Operators; +import org.apache.parquet.hadoop.metadata.ColumnPath; +import org.apache.parquet.hadoop.util.SerializationUtil; +import org.apache.parquet.io.api.Binary; import org.apache.pig.Expression; import org.apache.pig.LoadFunc; import org.apache.pig.LoadMetadata; +import org.apache.pig.LoadPredicatePushdown; import org.apache.pig.LoadPushDown; import org.apache.pig.ResourceSchema; import org.apache.pig.ResourceStatistics; @@ -57,6 +68,8 @@ import org.apache.pig.impl.util.UDFContext; import org.apache.pig.parser.ParserException; +import static org.apache.pig.Expression.*; + import org.apache.parquet.Log; import org.apache.parquet.hadoop.ParquetInputFormat; import org.apache.parquet.hadoop.metadata.GlobalMetaData; @@ -70,9 +83,12 @@ * @author Julien Le Dem * */ -public class ParquetLoader extends LoadFunc implements LoadMetadata, LoadPushDown { +public class ParquetLoader extends LoadFunc implements LoadMetadata, LoadPushDown, LoadPredicatePushdown { private static final Log LOG = Log.getLog(ParquetLoader.class); + public static final String ENABLE_PREDICATE_FILTER_PUSHDOWN = "parquet.pig.predicate.pushdown.enable"; + private static final boolean DEFAULT_PREDICATE_PUSHDOWN_ENABLED = false; + // Using a weak hash map will ensure that the cache will be gc'ed when there is memory pressure static final Map>> inputFormatCache = new WeakHashMap>>(); @@ -169,6 +185,7 @@ private void setInput(String location, Job job) throws IOException { getConfiguration(job).set(PARQUET_PIG_SCHEMA, pigSchemaToString(schema)); getConfiguration(job).set(PARQUET_PIG_REQUIRED_FIELDS, serializeRequiredFieldList(requiredFieldList)); getConfiguration(job).set(PARQUET_COLUMN_INDEX_ACCESS, Boolean.toString(columnIndexAccess)); + SerializationUtil.writeObjectToConfAsBase64(ParquetInputFormat.FILTER_PREDICATE, getFromUDFContext(ParquetInputFormat.FILTER_PREDICATE), getConfiguration(job)); } @Override @@ -380,4 +397,133 @@ private Schema getSchemaFromRequiredFieldList(Schema schema, List return s; } + @Override + public List getPredicateFields(String s, Job job) throws IOException { + if(!job.getConfiguration().getBoolean(ENABLE_PREDICATE_FILTER_PUSHDOWN, DEFAULT_PREDICATE_PUSHDOWN_ENABLED)) { + return null; + } + + List fields = new ArrayList(); + + for(FieldSchema field : schema.getFields()) { + switch(field.type) { + case DataType.BOOLEAN: + case DataType.INTEGER: + case DataType.LONG: + case DataType.FLOAT: + case DataType.DOUBLE: + case DataType.CHARARRAY: + fields.add(field.alias); + break; + default: + // Skip BYTEARRAY, TUPLE, MAP, BAG, DATETIME, BIGINTEGER, BIGDECIMAL + break; + } + } + + return fields; + } + + @Override + public List getSupportedExpressionTypes() { + Expression.OpType supportedTypes [] = { + Expression.OpType.OP_EQ, + Expression.OpType.OP_GT, + Expression.OpType.OP_GE, + Expression.OpType.OP_LT, + Expression.OpType.OP_LE, + Expression.OpType.OP_AND, + Expression.OpType.OP_OR, + }; + + return Arrays.asList(supportedTypes); + } + + @Override + public void setPushdownPredicate(Expression e) throws IOException { + LOG.info("Pig pushdown expression:" + e); + + FilterPredicate pred = buildFilter(e); + LOG.info("Pig pushdown expression:" + e); + + storeInUDFContext(ParquetInputFormat.FILTER_PREDICATE, pred); + } + + private FilterPredicate buildFilter(Expression e) { + if (e instanceof BinaryExpression) { + Expression lhs = ((BinaryExpression) e).getLhs(); + Expression rhs = ((BinaryExpression) e).getRhs(); + + switch (e.getOpType()) { + case OP_EQ: return eq(getColumn((Column) lhs), getValue(rhs)); + case OP_GT: return gt(getColumn((Column) lhs), getValue(rhs)); + case OP_GE: return gtEq(getColumn((Column) lhs), getValue(rhs)); + case OP_LT: return lt(getColumn((Column) lhs), getValue(rhs)); + case OP_LE: return ltEq(getColumn((Column) lhs), getValue(rhs)); + case OP_AND: return and(buildFilter(lhs), buildFilter(rhs)); + case OP_OR: return or(buildFilter(lhs), buildFilter(rhs)); + } + } + + return null; + } + + private ColumnWrapper getColumn(Column c) { + String column = c.getName(); + + try { + FieldSchema f = schema.getField(column); + + switch (f.type) { + case DataType.BOOLEAN: return new ColumnWrapper(binaryColumn(f.alias)); + case DataType.INTEGER: return new ColumnWrapper(intColumn(f.alias)); + case DataType.LONG: return new ColumnWrapper(longColumn(f.alias)); + case DataType.FLOAT: return new ColumnWrapper(floatColumn(f.alias)); + case DataType.DOUBLE: return new ColumnWrapper(doubleColumn(f.alias)); + case DataType.CHARARRAY: return new ColumnWrapper(binaryColumn(f.alias)); + } + + } catch (FrontendException e) { + throw new RuntimeException("Error processing pushdown for column:" + c , e); + } + + return null; + } + + private static class ColumnWrapper> extends Operators.Column implements Operators.SupportsLtGt, Serializable { + + Operators.Column wrapped; + + public ColumnWrapper(Operators.Column wrapped) { + super(wrapped.getColumnPath(), wrapped.getColumnType()); + this.wrapped = wrapped; + } + + @Override + public Class getColumnType() { + return wrapped.getColumnType(); + } + + @Override + public ColumnPath getColumnPath() { + return wrapped.getColumnPath(); + } + } + + private Comparable getValue(Expression expr) { + switch(expr.getOpType()) { + case TERM_COL: + //return ((Column) expr).getName(); + case TERM_CONST: + Comparable value = (Comparable) ((Const) expr).getValue(); + + if(value instanceof String) { + value = Binary.fromString((String)value); + } + return value; + default: + throw new RuntimeException("Unsupported expression type: " + expr.getOpType() + " in " + expr); + } + } + } diff --git a/parquet-pig/src/test/java/org/apache/parquet/pig/TestParquetLoader.java b/parquet-pig/src/test/java/org/apache/parquet/pig/TestParquetLoader.java index 6f11538d85..c151b929a4 100644 --- a/parquet-pig/src/test/java/org/apache/parquet/pig/TestParquetLoader.java +++ b/parquet-pig/src/test/java/org/apache/parquet/pig/TestParquetLoader.java @@ -22,6 +22,8 @@ import java.util.Arrays; import java.util.List; import java.util.Properties; + +import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.mapreduce.Job; import org.apache.pig.ExecType; import org.apache.pig.LoadPushDown.RequiredField; @@ -175,11 +177,11 @@ public void testReqestedSchemaColumnPruning() throws Exception { for (int i = 0; i < rows; i++) { list.add(Storage.tuple(i, "a"+i, i*2)); } - data.set("in", "i:int, a:chararray, b:int", list ); + data.set("in", "i:int, a:chararray, b:int", list); pigServer.setBatchOn(); pigServer.registerQuery("A = LOAD 'in' USING mock.Storage();"); pigServer.deleteFile(out); - pigServer.registerQuery("Store A into '"+out+"' using " + ParquetStorer.class.getName()+"();"); + pigServer.registerQuery("Store A into '" + out + "' using " + ParquetStorer.class.getName() + "();"); pigServer.executeBatch(); //Test Null Padding at the end @@ -212,7 +214,7 @@ public void testTypePersuasion() throws Exception { for (int i = 0; i < rows; i++) { list.add(Storage.tuple(i, (long)i, (float)i, (double)i, Integer.toString(i), Boolean.TRUE)); } - data.set("in", "i:int, l:long, f:float, d:double, s:chararray, b:boolean", list ); + data.set("in", "i:int, l:long, f:float, d:double, s:chararray, b:boolean", list); pigServer.setBatchOn(); pigServer.registerQuery("A = LOAD 'in' USING mock.Storage();"); pigServer.deleteFile(out); @@ -268,11 +270,11 @@ public void testColumnIndexAccess() throws Exception { pigServer.setBatchOn(); pigServer.registerQuery("A = LOAD 'in' USING mock.Storage();"); pigServer.deleteFile(out); - pigServer.registerQuery("Store A into '"+out+"' using " + ParquetStorer.class.getName()+"();"); + pigServer.registerQuery("Store A into '" + out + "' using " + ParquetStorer.class.getName() + "();"); pigServer.executeBatch(); //Test Null Padding at the end - pigServer.registerQuery("B = LOAD '" + out + "' using " + ParquetLoader.class.getName()+"('n1:int, n2:double, n3:long, n4:chararray', 'true');"); + pigServer.registerQuery("B = LOAD '" + out + "' using " + ParquetLoader.class.getName() + "('n1:int, n2:double, n3:long, n4:chararray', 'true');"); pigServer.registerQuery("STORE B into 'out' using mock.Storage();"); pigServer.executeBatch(); @@ -285,7 +287,7 @@ public void testColumnIndexAccess() throws Exception { assertEquals(4, t.size()); assertEquals(i, t.get(0)); - assertEquals(i*1.0, t.get(1)); + assertEquals(i * 1.0, t.get(1)); assertEquals(i*2L, t.get(2)); assertEquals("v"+i, t.get(3)); } @@ -306,10 +308,10 @@ public void testColumnIndexAccessProjection() throws Exception { pigServer.setBatchOn(); pigServer.registerQuery("A = LOAD 'in' USING mock.Storage();"); pigServer.deleteFile(out); - pigServer.registerQuery("Store A into '"+out+"' using " + ParquetStorer.class.getName()+"();"); + pigServer.registerQuery("Store A into '" + out + "' using " + ParquetStorer.class.getName() + "();"); pigServer.executeBatch(); - pigServer.registerQuery("B = LOAD '" + out + "' using " + ParquetLoader.class.getName()+"('n1:int, n2:double, n3:long, n4:chararray', 'true');"); + pigServer.registerQuery("B = LOAD '" + out + "' using " + ParquetLoader.class.getName() + "('n1:int, n2:double, n3:long, n4:chararray', 'true');"); pigServer.registerQuery("C = foreach B generate n1, n3;"); pigServer.registerQuery("STORE C into 'out' using mock.Storage();"); pigServer.executeBatch(); @@ -325,10 +327,37 @@ public void testColumnIndexAccessProjection() throws Exception { assertEquals(i, t.get(0)); assertEquals(i*2L, t.get(1)); } - } - + } + @Test - public void testRead() { - + public void testPredicatePushdown() throws Exception { + Configuration conf = new Configuration(); + conf.setBoolean(ParquetLoader.ENABLE_PREDICATE_FILTER_PUSHDOWN, true); + + PigServer pigServer = new PigServer(ExecType.LOCAL); + pigServer.setValidateEachStatement(true); + + String out = "target/out"; + int rows = 10; + Data data = Storage.resetData(pigServer); + List list = new ArrayList(); + for (int i = 0; i < rows; i++) { + list.add(Storage.tuple(i, i*1.0, i*2L, "v"+i)); + } + data.set("in", "c1:int, c2:double, c3:long, c4:chararray", list); + pigServer.setBatchOn(); + pigServer.registerQuery("A = LOAD 'in' USING mock.Storage();"); + pigServer.deleteFile(out); + pigServer.registerQuery("Store A into '"+out+"' using " + ParquetStorer.class.getName()+"();"); + pigServer.executeBatch(); + + pigServer.registerQuery("B = LOAD '" + out + "' using " + ParquetLoader.class.getName()+"('c1:int, c2:double, c3:long, c4:chararray');"); + pigServer.registerQuery("C = FILTER B by c1 == 1 or c1 == 5;"); + pigServer.registerQuery("STORE C into 'out' using mock.Storage();"); + pigServer.executeBatch(); + + List actualList = data.get("out"); + + assertEquals(2, actualList.size()); } } diff --git a/pom.xml b/pom.xml index c769ad3696..b2a760cd61 100644 --- a/pom.xml +++ b/pom.xml @@ -87,7 +87,7 @@ 2.10 false - 0.11.1 + 0.14.0 0.7.0 6.5.7 @@ -489,7 +489,7 @@ true 2.3.0 - 0.13.0 + 0.14.0 h2 From 7b019a6e2f7e6e256e669930dd3b54cca44b52c3 Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Tue, 1 Dec 2015 12:18:52 -0800 Subject: [PATCH 2/8] update tests and logging --- .../src/main/java/org/apache/parquet/pig/ParquetLoader.java | 4 ++-- .../test/java/org/apache/parquet/pig/TestParquetLoader.java | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java index 69227afd25..405abeec8c 100644 --- a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java +++ b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java @@ -441,10 +441,10 @@ public List getSupportedExpressionTypes() { @Override public void setPushdownPredicate(Expression e) throws IOException { - LOG.info("Pig pushdown expression:" + e); + LOG.info("Pig pushdown expression: " + e); FilterPredicate pred = buildFilter(e); - LOG.info("Pig pushdown expression:" + e); + LOG.info("Parquet filter predicate expression: " + pred); storeInUDFContext(ParquetInputFormat.FILTER_PREDICATE, pred); } diff --git a/parquet-pig/src/test/java/org/apache/parquet/pig/TestParquetLoader.java b/parquet-pig/src/test/java/org/apache/parquet/pig/TestParquetLoader.java index c151b929a4..1362254689 100644 --- a/parquet-pig/src/test/java/org/apache/parquet/pig/TestParquetLoader.java +++ b/parquet-pig/src/test/java/org/apache/parquet/pig/TestParquetLoader.java @@ -334,7 +334,7 @@ public void testPredicatePushdown() throws Exception { Configuration conf = new Configuration(); conf.setBoolean(ParquetLoader.ENABLE_PREDICATE_FILTER_PUSHDOWN, true); - PigServer pigServer = new PigServer(ExecType.LOCAL); + PigServer pigServer = new PigServer(ExecType.LOCAL, conf); pigServer.setValidateEachStatement(true); String out = "target/out"; From 266684962afcb74513e5c7fc4dc5c4334082b1f7 Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Tue, 1 Dec 2015 13:28:41 -0800 Subject: [PATCH 3/8] Fixed test to check the actual number of materialized rows from the reader --- .../apache/parquet/pig/TestParquetLoader.java | 16 ++++++++++------ 1 file changed, 10 insertions(+), 6 deletions(-) diff --git a/parquet-pig/src/test/java/org/apache/parquet/pig/TestParquetLoader.java b/parquet-pig/src/test/java/org/apache/parquet/pig/TestParquetLoader.java index 1362254689..8e57424498 100644 --- a/parquet-pig/src/test/java/org/apache/parquet/pig/TestParquetLoader.java +++ b/parquet-pig/src/test/java/org/apache/parquet/pig/TestParquetLoader.java @@ -29,12 +29,14 @@ import org.apache.pig.LoadPushDown.RequiredField; import org.apache.pig.LoadPushDown.RequiredFieldList; import org.apache.pig.PigServer; +import org.apache.pig.backend.executionengine.ExecJob; import org.apache.pig.builtin.mock.Storage; import org.apache.pig.builtin.mock.Storage.Data; import org.apache.pig.data.DataType; import static org.apache.pig.data.DataType.*; import org.apache.pig.data.Tuple; import org.apache.pig.impl.logicalLayer.FrontendException; +import org.apache.pig.tools.pigstats.JobStats; import org.junit.Assert; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertEquals; @@ -338,6 +340,7 @@ public void testPredicatePushdown() throws Exception { pigServer.setValidateEachStatement(true); String out = "target/out"; + String out2 = "target/out2"; int rows = 10; Data data = Storage.resetData(pigServer); List list = new ArrayList(); @@ -348,16 +351,17 @@ public void testPredicatePushdown() throws Exception { pigServer.setBatchOn(); pigServer.registerQuery("A = LOAD 'in' USING mock.Storage();"); pigServer.deleteFile(out); - pigServer.registerQuery("Store A into '"+out+"' using " + ParquetStorer.class.getName()+"();"); + pigServer.registerQuery("Store A into '" + out + "' using " + ParquetStorer.class.getName() + "();"); pigServer.executeBatch(); - pigServer.registerQuery("B = LOAD '" + out + "' using " + ParquetLoader.class.getName()+"('c1:int, c2:double, c3:long, c4:chararray');"); + pigServer.deleteFile(out2); + pigServer.registerQuery("B = LOAD '" + out + "' using " + ParquetLoader.class.getName() + "('c1:int, c2:double, c3:long, c4:chararray');"); pigServer.registerQuery("C = FILTER B by c1 == 1 or c1 == 5;"); - pigServer.registerQuery("STORE C into 'out' using mock.Storage();"); - pigServer.executeBatch(); + pigServer.registerQuery("STORE C into '" + out2 +"' using mock.Storage();"); + List jobs = pigServer.executeBatch(); - List actualList = data.get("out"); + long recordsRead = jobs.get(0).getStatistics().getInputStats().get(0).getNumberRecords(); - assertEquals(2, actualList.size()); + assertEquals(2, recordsRead); } } From a39fdffcf171b0806c71cb9d0d78b05dac609d38 Mon Sep 17 00:00:00 2001 From: Ryan Blue Date: Mon, 7 Dec 2015 16:32:42 -0800 Subject: [PATCH 4/8] WIP: Handle a few error cases in Pig predicate pushdown. --- .../org/apache/parquet/pig/ParquetLoader.java | 115 ++++++++++-------- 1 file changed, 65 insertions(+), 50 deletions(-) diff --git a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java index 405abeec8c..bebed3637f 100644 --- a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java +++ b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java @@ -185,7 +185,8 @@ private void setInput(String location, Job job) throws IOException { getConfiguration(job).set(PARQUET_PIG_SCHEMA, pigSchemaToString(schema)); getConfiguration(job).set(PARQUET_PIG_REQUIRED_FIELDS, serializeRequiredFieldList(requiredFieldList)); getConfiguration(job).set(PARQUET_COLUMN_INDEX_ACCESS, Boolean.toString(columnIndexAccess)); - SerializationUtil.writeObjectToConfAsBase64(ParquetInputFormat.FILTER_PREDICATE, getFromUDFContext(ParquetInputFormat.FILTER_PREDICATE), getConfiguration(job)); + SerializationUtil.writeObjectToConfAsBase64(ParquetInputFormat.FILTER_PREDICATE, + getFromUDFContext(ParquetInputFormat.FILTER_PREDICATE), getConfiguration(job)); } @Override @@ -453,77 +454,91 @@ private FilterPredicate buildFilter(Expression e) { if (e instanceof BinaryExpression) { Expression lhs = ((BinaryExpression) e).getLhs(); Expression rhs = ((BinaryExpression) e).getRhs(); + OpType op = e.getOpType(); + + FilterPredicate lfp; + FilterPredicate rfp; + switch (op) { + case OP_AND: + lfp = buildFilter(lhs); + rfp = buildFilter(rhs); + if (lfp == null || rfp == null) { + return null; + } + return and(lfp, rfp); + case OP_OR: + lfp = buildFilter(lhs); + rfp = buildFilter(rhs); + if (lfp == null || rfp == null) { + return null; + } + return or(lfp, rfp); + } - switch (e.getOpType()) { - case OP_EQ: return eq(getColumn((Column) lhs), getValue(rhs)); - case OP_GT: return gt(getColumn((Column) lhs), getValue(rhs)); - case OP_GE: return gtEq(getColumn((Column) lhs), getValue(rhs)); - case OP_LT: return lt(getColumn((Column) lhs), getValue(rhs)); - case OP_LE: return ltEq(getColumn((Column) lhs), getValue(rhs)); - case OP_AND: return and(buildFilter(lhs), buildFilter(rhs)); - case OP_OR: return or(buildFilter(lhs), buildFilter(rhs)); + if (lhs instanceof Column && rhs instanceof Const) { + return buildFilter(op, (Column) lhs, (Const) rhs); + } else if (lhs instanceof Const && rhs instanceof Column) { + return buildFilter(op, (Column) rhs, (Const) lhs); } } return null; } - private ColumnWrapper getColumn(Column c) { - String column = c.getName(); - + private FilterPredicate buildFilter(OpType op, Column col, Const value) { + String name = col.getName(); try { - FieldSchema f = schema.getField(column); - + FieldSchema f = schema.getField(name); switch (f.type) { - case DataType.BOOLEAN: return new ColumnWrapper(binaryColumn(f.alias)); - case DataType.INTEGER: return new ColumnWrapper(intColumn(f.alias)); - case DataType.LONG: return new ColumnWrapper(longColumn(f.alias)); - case DataType.FLOAT: return new ColumnWrapper(floatColumn(f.alias)); - case DataType.DOUBLE: return new ColumnWrapper(doubleColumn(f.alias)); - case DataType.CHARARRAY: return new ColumnWrapper(binaryColumn(f.alias)); +// case DataType.BOOLEAN: +// Operators.BooleanColumn col = booleanColumn(name); +// return op(op, col, value); + case DataType.INTEGER: + Operators.IntColumn intCol = intColumn(name); + return op(op, intCol, value); + case DataType.LONG: + Operators.LongColumn longCol = longColumn(name); + return op(op, longCol, value); + case DataType.FLOAT: + Operators.FloatColumn floatCol = floatColumn(name); + return op(op, floatCol, value); + case DataType.DOUBLE: + Operators.DoubleColumn doubleCol = doubleColumn(name); + return op(op, doubleCol, value); + case DataType.CHARARRAY: + Operators.BinaryColumn binaryCol = binaryColumn(name); + return op(op, binaryCol, value); + } } catch (FrontendException e) { - throw new RuntimeException("Error processing pushdown for column:" + c , e); + throw new RuntimeException("Error processing pushdown for column:" + col, e); } return null; } - private static class ColumnWrapper> extends Operators.Column implements Operators.SupportsLtGt, Serializable { - - Operators.Column wrapped; - - public ColumnWrapper(Operators.Column wrapped) { - super(wrapped.getColumnPath(), wrapped.getColumnType()); - this.wrapped = wrapped; - } - - @Override - public Class getColumnType() { - return wrapped.getColumnType(); - } - - @Override - public ColumnPath getColumnPath() { - return wrapped.getColumnPath(); + private , COL extends Operators.Column & Operators.SupportsLtGt> + FilterPredicate op(Expression.OpType op, COL col, Expression valueExpr) { + C value = getValue(valueExpr, col.getColumnType()); + switch (op) { + case OP_EQ: return eq(col, value); + case OP_GT: return gt(col, value); + case OP_GE: return gtEq(col, value); + case OP_LT: return lt(col, value); + case OP_LE: return ltEq(col, value); } + return null; } - private Comparable getValue(Expression expr) { - switch(expr.getOpType()) { - case TERM_COL: - //return ((Column) expr).getName(); - case TERM_CONST: - Comparable value = (Comparable) ((Const) expr).getValue(); + private > C getValue(Expression expr, Class type) { + Comparable value = (Comparable) ((Const) expr).getValue(); - if(value instanceof String) { - value = Binary.fromString((String)value); - } - return value; - default: - throw new RuntimeException("Unsupported expression type: " + expr.getOpType() + " in " + expr); + if (value instanceof String) { + value = Binary.fromString((String) value); } + + return type.cast(value); } } From f1ef73eb9ebab30f96e0f08b6ecb38b1c47f1097 Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Wed, 16 Dec 2015 15:31:56 -0800 Subject: [PATCH 5/8] Fixed binary type and storing filter predicate --- .../org/apache/parquet/pig/ParquetLoader.java | 19 +++++++++++-------- 1 file changed, 11 insertions(+), 8 deletions(-) diff --git a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java index 405abeec8c..9b3840b4ad 100644 --- a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java +++ b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java @@ -185,7 +185,11 @@ private void setInput(String location, Job job) throws IOException { getConfiguration(job).set(PARQUET_PIG_SCHEMA, pigSchemaToString(schema)); getConfiguration(job).set(PARQUET_PIG_REQUIRED_FIELDS, serializeRequiredFieldList(requiredFieldList)); getConfiguration(job).set(PARQUET_COLUMN_INDEX_ACCESS, Boolean.toString(columnIndexAccess)); - SerializationUtil.writeObjectToConfAsBase64(ParquetInputFormat.FILTER_PREDICATE, getFromUDFContext(ParquetInputFormat.FILTER_PREDICATE), getConfiguration(job)); + + FilterPredicate filterPredicate = (FilterPredicate) getFromUDFContext(ParquetInputFormat.FILTER_PREDICATE); + if(filterPredicate != null) { + ParquetInputFormat.setFilterPredicate(getConfiguration(job), filterPredicate); + } } @Override @@ -462,10 +466,11 @@ private FilterPredicate buildFilter(Expression e) { case OP_LE: return ltEq(getColumn((Column) lhs), getValue(rhs)); case OP_AND: return and(buildFilter(lhs), buildFilter(rhs)); case OP_OR: return or(buildFilter(lhs), buildFilter(rhs)); + default: throw new RuntimeException("Unsupported operation type: " + e.getOpType()); } + } else { + throw new RuntimeException("Unsupported expression type: " + e); } - - return null; } private ColumnWrapper getColumn(Column c) { @@ -475,19 +480,17 @@ private ColumnWrapper getColumn(Column c) { FieldSchema f = schema.getField(column); switch (f.type) { - case DataType.BOOLEAN: return new ColumnWrapper(binaryColumn(f.alias)); + case DataType.BOOLEAN: return new ColumnWrapper(booleanColumn(f.alias)); case DataType.INTEGER: return new ColumnWrapper(intColumn(f.alias)); case DataType.LONG: return new ColumnWrapper(longColumn(f.alias)); case DataType.FLOAT: return new ColumnWrapper(floatColumn(f.alias)); case DataType.DOUBLE: return new ColumnWrapper(doubleColumn(f.alias)); case DataType.CHARARRAY: return new ColumnWrapper(binaryColumn(f.alias)); + default: throw new RuntimeException("Error processing pushdown for column:" + c); } - } catch (FrontendException e) { throw new RuntimeException("Error processing pushdown for column:" + c , e); } - - return null; } private static class ColumnWrapper> extends Operators.Column implements Operators.SupportsLtGt, Serializable { @@ -513,7 +516,7 @@ public ColumnPath getColumnPath() { private Comparable getValue(Expression expr) { switch(expr.getOpType()) { case TERM_COL: - //return ((Column) expr).getName(); + throw new RuntimeException("Column comparison is unsupported: " + expr); case TERM_CONST: Comparable value = (Comparable) ((Const) expr).getValue(); From 388099b44be4921d960f7c8c5bbcf20fd40f5f5d Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Mon, 4 Jan 2016 12:54:39 -0800 Subject: [PATCH 6/8] Cleaning up imports --- .../org/apache/parquet/pig/ParquetLoader.java | 33 ++++++++++--------- 1 file changed, 18 insertions(+), 15 deletions(-) diff --git a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java index 1d82583bf8..151e0f2aa3 100644 --- a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java +++ b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java @@ -33,7 +33,6 @@ import static org.apache.parquet.filter2.predicate.FilterApi.*; import java.io.IOException; -import java.io.Serializable; import java.lang.ref.Reference; import java.lang.ref.SoftReference; import java.util.ArrayList; @@ -49,8 +48,6 @@ import org.apache.hadoop.mapreduce.TaskAttemptContext; import org.apache.parquet.filter2.predicate.FilterPredicate; import org.apache.parquet.filter2.predicate.Operators; -import org.apache.parquet.hadoop.metadata.ColumnPath; -import org.apache.parquet.hadoop.util.SerializationUtil; import org.apache.parquet.io.api.Binary; import org.apache.pig.Expression; import org.apache.pig.LoadFunc; @@ -68,7 +65,10 @@ import org.apache.pig.impl.util.UDFContext; import org.apache.pig.parser.ParserException; -import static org.apache.pig.Expression.*; +import static org.apache.pig.Expression.BinaryExpression; +import static org.apache.pig.Expression.Column; +import static org.apache.pig.Expression.Const; +import static org.apache.pig.Expression.OpType; import org.apache.parquet.Log; import org.apache.parquet.hadoop.ParquetInputFormat; @@ -430,14 +430,14 @@ public List getPredicateFields(String s, Job job) throws IOException { @Override public List getSupportedExpressionTypes() { - Expression.OpType supportedTypes [] = { - Expression.OpType.OP_EQ, - Expression.OpType.OP_GT, - Expression.OpType.OP_GE, - Expression.OpType.OP_LT, - Expression.OpType.OP_LE, - Expression.OpType.OP_AND, - Expression.OpType.OP_OR, + OpType supportedTypes [] = { + OpType.OP_EQ, + OpType.OP_GT, + OpType.OP_GE, + OpType.OP_LT, + OpType.OP_LE, + OpType.OP_AND, + OpType.OP_OR }; return Arrays.asList(supportedTypes); @@ -493,9 +493,12 @@ private FilterPredicate buildFilter(OpType op, Column col, Const value) { try { FieldSchema f = schema.getField(name); switch (f.type) { -// case DataType.BOOLEAN: -// Operators.BooleanColumn col = booleanColumn(name); -// return op(op, col, value); + case DataType.BOOLEAN: + Operators.BooleanColumn boolCol = booleanColumn(name); + switch(op) { + case OP_EQ: return eq(boolCol, getValue(value, boolCol.getColumnType())); + case OP_NE: return notEq(boolCol, getValue(value, boolCol.getColumnType())); + } case DataType.INTEGER: Operators.IntColumn intCol = intColumn(name); return op(op, intCol, value); From 54e23a6ed0a82479977d358e0f652c9604b8d0f2 Mon Sep 17 00:00:00 2001 From: Ryan Blue Date: Fri, 19 Feb 2016 16:02:04 -0800 Subject: [PATCH 7/8] PARQUET-397: Update Pig PPD to throw for bad expressions. Previously, the expresison builder would return null for unsupported expressions, but it is unclear whether that behavior is correct. Instead, if the expression can't be converted to a filter, this will throw an exception so it can be rewritten. This also implements "in", "between", "!=", and "not" expressions. --- .../org/apache/parquet/pig/ParquetLoader.java | 62 ++++++++++++------- 1 file changed, 40 insertions(+), 22 deletions(-) diff --git a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java index 151e0f2aa3..5f4be6ec1a 100644 --- a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java +++ b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java @@ -47,9 +47,13 @@ import org.apache.hadoop.mapreduce.RecordReader; import org.apache.hadoop.mapreduce.TaskAttemptContext; import org.apache.parquet.filter2.predicate.FilterPredicate; +import org.apache.parquet.filter2.predicate.LogicalInverseRewriter; import org.apache.parquet.filter2.predicate.Operators; import org.apache.parquet.io.api.Binary; import org.apache.pig.Expression; +import org.apache.pig.Expression.BetweenExpression; +import org.apache.pig.Expression.InExpression; +import org.apache.pig.Expression.UnaryExpression; import org.apache.pig.LoadFunc; import org.apache.pig.LoadMetadata; import org.apache.pig.LoadPredicatePushdown; @@ -432,12 +436,16 @@ public List getPredicateFields(String s, Job job) throws IOException { public List getSupportedExpressionTypes() { OpType supportedTypes [] = { OpType.OP_EQ, + OpType.OP_NE, OpType.OP_GT, OpType.OP_GE, OpType.OP_LT, OpType.OP_LE, OpType.OP_AND, - OpType.OP_OR + OpType.OP_OR, + OpType.OP_BETWEEN, + OpType.OP_IN, + OpType.OP_NOT }; return Arrays.asList(supportedTypes); @@ -454,28 +462,33 @@ public void setPushdownPredicate(Expression e) throws IOException { } private FilterPredicate buildFilter(Expression e) { + OpType op = e.getOpType(); + if (e instanceof BinaryExpression) { Expression lhs = ((BinaryExpression) e).getLhs(); Expression rhs = ((BinaryExpression) e).getRhs(); - OpType op = e.getOpType(); - FilterPredicate lfp; - FilterPredicate rfp; switch (op) { case OP_AND: - lfp = buildFilter(lhs); - rfp = buildFilter(rhs); - if (lfp == null || rfp == null) { - return null; - } - return and(lfp, rfp); + return and(buildFilter(lhs), buildFilter(rhs)); case OP_OR: - lfp = buildFilter(lhs); - rfp = buildFilter(rhs); - if (lfp == null || rfp == null) { - return null; + return or(buildFilter(lhs), buildFilter(rhs)); + case OP_BETWEEN: + BetweenExpression between = (BetweenExpression) rhs; + return and( + buildFilter(OpType.OP_GE, (Column) lhs, (Const) between.getLower()), + buildFilter(OpType.OP_LE, (Column) lhs, (Const) between.getUpper())); + case OP_IN: + FilterPredicate current = null; + for (Object value : ((InExpression) rhs).getValues()) { + FilterPredicate next = buildFilter(OpType.OP_EQ, (Column) lhs, (Const) value); + if (current != null) { + current = or(current, next); + } else { + current = next; + } } - return or(lfp, rfp); + return current; } if (lhs instanceof Column && rhs instanceof Const) { @@ -483,9 +496,12 @@ private FilterPredicate buildFilter(Expression e) { } else if (lhs instanceof Const && rhs instanceof Column) { return buildFilter(op, (Column) rhs, (Const) lhs); } + } else if (e instanceof UnaryExpression && op == OpType.OP_NOT) { + return LogicalInverseRewriter.rewrite( + not(buildFilter(((UnaryExpression) e).getExpression()))); } - return null; + throw new RuntimeException("Could not build filter for expression: " + e); } private FilterPredicate buildFilter(OpType op, Column col, Const value) { @@ -498,6 +514,8 @@ private FilterPredicate buildFilter(OpType op, Column col, Const value) { switch(op) { case OP_EQ: return eq(boolCol, getValue(value, boolCol.getColumnType())); case OP_NE: return notEq(boolCol, getValue(value, boolCol.getColumnType())); + default: throw new RuntimeException( + "Operation " + op + " not supported for boolean column: " + name); } case DataType.INTEGER: Operators.IntColumn intCol = intColumn(name); @@ -514,20 +532,20 @@ private FilterPredicate buildFilter(OpType op, Column col, Const value) { case DataType.CHARARRAY: Operators.BinaryColumn binaryCol = binaryColumn(name); return op(op, binaryCol, value); - + default: + throw new RuntimeException("Unsupported type " + f.type + " for field: " + name); } } catch (FrontendException e) { throw new RuntimeException("Error processing pushdown for column:" + col, e); } - - return null; } private , COL extends Operators.Column & Operators.SupportsLtGt> - FilterPredicate op(Expression.OpType op, COL col, Expression valueExpr) { + FilterPredicate op(Expression.OpType op, COL col, Const valueExpr) { C value = getValue(valueExpr, col.getColumnType()); switch (op) { case OP_EQ: return eq(col, value); + case OP_NE: return notEq(col, value); case OP_GT: return gt(col, value); case OP_GE: return gtEq(col, value); case OP_LT: return lt(col, value); @@ -536,8 +554,8 @@ FilterPredicate op(Expression.OpType op, COL col, Expression valueExpr) { return null; } - private > C getValue(Expression expr, Class type) { - Comparable value = (Comparable) ((Const) expr).getValue(); + private > C getValue(Const valueExpr, Class type) { + Object value = valueExpr.getValue(); if (value instanceof String) { value = Binary.fromString((String) value); From c7a9b02a8c38b973c7acc81ee29d4a8cde72ec52 Mon Sep 17 00:00:00 2001 From: Ryan Blue Date: Fri, 26 Feb 2016 10:26:55 -0800 Subject: [PATCH 8/8] PARQUET-397: Address review comments. --- .../main/java/org/apache/parquet/pig/ParquetLoader.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java index 5f4be6ec1a..66777345a2 100644 --- a/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java +++ b/parquet-pig/src/main/java/org/apache/parquet/pig/ParquetLoader.java @@ -443,8 +443,8 @@ public List getSupportedExpressionTypes() { OpType.OP_LE, OpType.OP_AND, OpType.OP_OR, - OpType.OP_BETWEEN, - OpType.OP_IN, + //OpType.OP_BETWEEN, // not implemented in Pig yet + //OpType.OP_IN, // not implemented in Pig yet OpType.OP_NOT }; @@ -540,7 +540,7 @@ private FilterPredicate buildFilter(OpType op, Column col, Const value) { } } - private , COL extends Operators.Column & Operators.SupportsLtGt> + private static , COL extends Operators.Column & Operators.SupportsLtGt> FilterPredicate op(Expression.OpType op, COL col, Const valueExpr) { C value = getValue(valueExpr, col.getColumnType()); switch (op) { @@ -554,7 +554,7 @@ FilterPredicate op(Expression.OpType op, COL col, Const valueExpr) { return null; } - private > C getValue(Const valueExpr, Class type) { + private static > C getValue(Const valueExpr, Class type) { Object value = valueExpr.getValue(); if (value instanceof String) {