From f8ca4478923f56cf0a91394ea5cba98ba8da5241 Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Tue, 1 Dec 2015 11:34:49 -0800 Subject: [PATCH 1/6] 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/6] 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/6] 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/6] 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/6] 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/6] 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);