Use PXF object store connector to read and write ORC-format data between Greengage DB and S3
The Optimized Row Columnar (ORC) file format is a column-oriented data storage format that offers improvements over text and RCFile formats in terms of both compression and performance. ORC is type-aware and is specifically designed for Hadoop workloads. ORC files store both the type of the data in the file and its encoding information.
In an ORC file, all columns within a single group of row data (also known as stripe) are stored together on disk. The columnar nature of the format enables column projection on read operations, helping avoid accessing unnecessary columns during a query. ORC also supports predicate pushdown with built-in indexes at the file, stripe, and row levels, moving the filter operation to the data loading phase. The connector supports ORC file versions v0 and v1.
This topic describes how to configure and use the PXF object store connector for reading and writing ORC data in an object store by using external tables and provides practical examples.
Data type mapping
To read and write ORC primitive data types in Greengage DB, map the ORC data values to the Greengage DB columns of the same type.
Read mapping
To read ORC scalar data types in Greengage DB, map ORC data values to Greengage DB columns of the same type using the following data type mapping.
| ORC physical type | ORC logical type | PXF / Greengage DB data type |
|---|---|---|
binary |
decimal |
NUMERIC |
binary |
timestamp |
TIMESTAMP |
byte[] |
string |
TEXT |
byte[] |
char |
BPCHAR |
byte[] |
varchar |
VARCHAR |
byte[] |
binary |
BYTEA |
Double |
float |
REAL |
Double |
double |
FLOAT8 |
Integer |
|
BOOLEAN |
Integer |
|
SMALLINT |
Integer |
|
SMALLINT |
Integer |
|
INTEGER |
Integer |
|
BIGINT |
Integer |
date |
DATE |
PXF supports the list ORC compound type for a subset of the ORC scalar types.
The map, union, and struct compound types are not supported.
| ORC compound type | PXF / Greengage DB data type |
|---|---|
array<string> |
TEXT[] |
array<char> |
BPCHAR[] |
array<varchar> |
VARCHAR[] |
array<binary> |
BYTEA[] |
array<float> |
REAL[] |
array<double> |
FLOAT8[] |
array<boolean> |
BOOLEAN[] |
array<tinyint> |
SMALLINT[] |
array<smallint> |
SMALLINT[] |
array<int> |
INTEGER[] |
array<bigint> |
BIGINT[] |
Write mapping
PXF uses the following data type mapping when writing ORC data.
| PXF / Greengage DB data type | ORC logical type | ORC physical type |
|---|---|---|
NUMERIC |
decimal |
binary |
TIMESTAMP (1) |
timestamp |
binary |
TIMESTAMPTZ |
timestamp with local timezone |
timestamp |
TEXT |
string |
byte[] |
BPCHAR |
char |
byte[] |
VARCHAR |
varchar |
byte[] |
BYTEA |
binary |
byte[] |
REAL |
float |
Double |
FLOAT8 |
double |
Double |
BOOLEAN |
|
Integer |
SMALLINT |
|
Integer |
SMALLINT |
|
Integer |
INTEGER |
|
Integer |
BIGINT |
|
Integer |
DATE |
date |
Integer |
UUID |
string |
byte[] |
-
The
pxf.orc.write.timezone.utcserver configuration property in the pxf-site.xml file defines how PXF writesTIMESTAMPvalues to the external data store. By default, aTIMESTAMPtype is written using the UTC timezone. To write aTIMESTAMPtype using the PXF JVM local timezone, set thepxf.orc.write.timezone.utcproperty tofalse.
Learn about configuring a PXF server in Configure a PXF server in PXF documentation.
PXF supports writing the list ORC compound type for one-dimensional arrays of all the ORC primitive types listed above.
The map, union, and struct compound types as well as user-provided schemas are not supported.
| PXF / Greengage DB data type | ORC compound type |
|---|---|
NUMERIC[] |
array<decimal> |
TIMESTAMP[] |
array<timestamp> |
TEXT[] |
array<string> |
BPCHAR[] |
array<char> |
VARCHAR[] |
array<varchar> |
BYTEA[] |
array<binary> |
REAL[] |
array<float> |
FLOAT8[] |
array<double> |
BOOLEAN[] |
array<boolean> |
SMALLINT[] |
array<tinyint> |
SMALLINT[] |
array<smallint> |
INTEGER[] |
array<int> |
BIGINT[] |
array<bigint> |
DATE[] |
array<date> |
Create an external table using the PXF protocol
To create a Greengage DB external table to read or write ORC data in an object store, use the following general syntax:
CREATE [READABLE | WRITABLE] EXTERNAL TABLE <table_name>
( <column_name> <data_type> [, ...] | LIKE <other_table> )
LOCATION ('pxf://<path_to_data>?PROFILE=<objstore>:orc[&<custom_option>=<value>[...]]')
FORMAT 'CUSTOM' (FORMATTER='pxfwritable_import'|'pxfwritable_export')
[DISTRIBUTED BY (<column_name> [, ... ] ) | DISTRIBUTED RANDOMLY];
| Keyword | Value |
|---|---|
<table_name> |
The name of the table to create |
<column_name> |
The name of the column to create |
<data_type> |
The data type of the column |
LIKE <other_table> |
Specifies a table from which the new external table automatically copies all column names, data types, and distribution policy |
<path_to_data> |
The path to the directory or file in an object store.
When the |
PROFILE=<objstore>:orc |
The profile is specified as the The following
|
FORMAT 'CUSTOM' |
The custom format with the built-in custom formatter functions for read ( |
DISTRIBUTED BY |
When loading data from a Greengage DB table into a writable external table, consider specifying the same distribution policy or column name on both tables. This will avoid extra motion of data between segments on the load operation. Learn more about table distribution in Distribution |
<custom_option> |
One of the custom options provided in the |
SERVER=<server_name> |
The named server configuration that PXF uses to access the data.
If the option is omitted, the |
IGNORE_MISSING_PATH |
The action to take when |
MAP_BY_POSITION |
Specifies whether PXF should map an ORC column to a Greengage DB column by position.
The default value is |
COMPRESSION_CODEC |
The compression codec to use when writing data: You must explicitly specify |
Numeric data overflow conditions
PXF uses the HiveDecimal class to write numeric ORC data, which limits both the precision and the scale of a numeric type to a maximum of 38.
When you define a NUMERIC column in an external table without specifying a precision or scale, PXF internally maps the column to DECIMAL(38, 18).
PXF handles the following precision overflow conditions:
-
A
NUMERICcolumn is defined in the external table, and the total digit count of a value exceeds the maximum supported precision of 38, for example,1234567890123456789012345678901234567890.12345, which has a total digit count of 45. -
A
NUMERIC(<precision>)column is defined with a<precision>value greater than 38, for exampleNUMERIC(55). -
A
NUMERICcolumn is defined in the external table, and the integer digit count of a value is greater than 20 (38-18), for example,123456789012345678901234567890.12345, which has an integer digit count of 30.
If you define a NUMERIC(<precision>, <scale>) column and the integer digit count of a value is greater than <precision>-<scale>, PXF returns an error.
For example, you define a NUMERIC(20,4) column and the value is 12345678901234567.12, whose integer digit count of 17 is greater than 20-4=16.
PXF can perform one of the following actions when detecting a numeric data overflow: round the value (the default), return an error, or ignore the overflow.
The pxf.orc.write.decimal.overflow property in the pxf-site.xml server configuration file specifies the action to take.
| Value | PXF action |
|---|---|
round |
The default behavior. When PXF encounters an overflow, it attempts to round the value to meet both precision and scale requirements before writing and reports an error if rounding fails. This may potentially leave an incomplete dataset in the external system |
error |
PXF reports an error when it encounters an overflow, and the transaction fails |
ignore |
PXF logs a warning and attempts to round the value to meet both precision and scale requirements; otherwise a |
Learn about configuring a PXF server in Configure a PXF server in PXF documentation.
Examples
These examples demonstrate how to configure and use the PXF object store connector for reading and writing ORC data in an object store by using external tables.
Configure the PXF S3 connector
To have PXF connect to an object store, you need to create the corresponding server configuration and then synchronize it to the Greengage DB cluster:
-
On the Greengage DB master host, log in as
gpadmin. -
Go to the $PXF_BASE/servers directory and create an S3 server configuration directory named s3. Depending on your object storage, copy the required server configuration file from $PXF_HOME/templates to $PXF_BASE/servers/s3. The example uses the configuration file based on the minio-site.xml template.
$ mkdir $PXF_BASE/servers/s3 $ cd $PXF_BASE/servers/s3 $ cp $PXF_HOME/templates/minio-site.xml .In the configuration file, provide the relevant object store connection details:
<?xml version="1.0" encoding="UTF-8"?> <configuration> <property> <name>fs.s3a.endpoint</name> <value>storage.example.com</value> </property> <property> <name>fs.s3a.access.key</name> <value>${ACCESS_KEY}</value> </property> <property> <name>fs.s3a.secret.key</name> <value>${SECRET_KEY}</value> </property> <property> <name>fs.s3a.fast.upload</name> <value>true</value> </property> <property> <name>fs.s3a.path.style.access</name> <value>true</value> </property> </configuration>NOTENotice that the credentials used to authenticate with the storage service are provided via the
ACCESS_KEYandSECRET_KEYenvironment variables.You can configure them as follows:
$ export ACCESS_KEY=<access_key> $ export SECRET_KEY=<secret_key>While it’s possible to set the credentials directly in the configuration file, it is recommended to use environment variables for security reasons.
-
Synchronize the server configuration to the Greengage DB cluster hosts:
$ pxf cluster sync
Create a readable external table
-
In the /tmp directory on your local machine, prepare a sample data set.
-
Create a file named customers.json having the following content:
{"id": 1,"name": "John Doe","ordered_items":["laptop", "monitor"]} {"id": 2,"name": "Jane Smith","ordered_items":["keyboard", "mouse", "pad"]} {"id": 3,"name": "Bob Brown","ordered_items":["headphones", "laptop"]} {"id": 4,"name": "Alice Green","ordered_items":["webcam", "microphone", "mouse"]}The following field names and data types are used in the file:
-
id—INT; -
name—TEXT; -
ordered_items—TEXT[].
-
-
Download the most recent version of the ORC Java tools JAR from Maven Central to the /tmp directory and convert customers.json to the customers.orc ORC file:
$ curl -O https://repo1.maven.org/maven2/org/apache/orc/orc-tools/2.2.2/orc-tools-2.2.2-uber.jar $ java -jar /tmp/orc-tools-2.2.2-uber.jar convert /tmp/customers.json \ --schema 'struct<id:int,name:string,ordered_items:array<string>>' \ -o /tmp/customers.orc
-
-
Add the generated ORC file to the customers bucket on the S3 host.
-
On the Greengage DB master host, create an external table that references the customers.orc file. In the
LOCATIONclause, specify the PXFs3:orcprofile and the server configuration. In theFORMATclause, specifypxfwritable_import, which is the built-in custom formatter function for read operations:CREATE EXTERNAL TABLE customers_r ( id INT, name TEXT, ordered_items TEXT[] ) LOCATION('pxf://customers/customers.orc?PROFILE=s3:orc&SERVER=s3') FORMAT 'CUSTOM' (FORMATTER='pxfwritable_import'); -
Query the created external table:
SELECT * FROM customers_r;The output should look as follows:
id | name | ordered_items ----+-------------+--------------------------- 1 | John Doe | {laptop,monitor} 2 | Jane Smith | {keyboard,mouse,pad} 3 | Bob Brown | {headphones,laptop} 4 | Alice Green | {webcam,microphone,mouse} (4 rows) -
Run a query to view the rows where the
ordered_itemscolumn includeslaptopormonitor:SELECT * FROM customers_r WHERE ordered_items && '{"laptop", "monitor"}';The output should look as follows:
id | name | ordered_items ----+-----------+--------------------- 1 | John Doe | {laptop,monitor} 3 | Bob Brown | {headphones,laptop} (2 rows) -
Run a query to view the rows where the first item in the
ordered_itemscolumn isheadphones:SELECT * FROM customers_r WHERE ordered_items[1] = 'headphones';The output should look as follows:
id | name | ordered_items ----+-----------+--------------------- 3 | Bob Brown | {headphones,laptop} (1 row)
Create a writable external table
-
On the Greengage DB master host, create the writable external table that stores data to the customers_w folder in the customers bucket on the S3 host. In the
LOCATIONclause, specify the PXFs3:orcprofile and the server configuration and setCOMPRESSION_CODECtononeto disable compression. In theFORMATclause, specifypxfwritable_export, which is the built-in custom formatter function for write operations:CREATE WRITABLE EXTERNAL TABLE customers_w ( id INT, name TEXT, ordered_items TEXT[] ) LOCATION ('pxf://customers/customers_w?PROFILE=s3:orc&SERVER=s3&COMPRESSION_CODEC=none') FORMAT 'CUSTOM' (FORMATTER='pxfwritable_export'); -
Insert some data into the
customers_wtable:INSERT INTO customers_w VALUES (1, 'Bob Brown', ARRAY['monitor', 'webcam']), (2, 'John Doe', ARRAY['keyboard', 'microphone']), (3, 'Alice Green', ARRAY['headphones', 'pad']), (4, 'Jane Smith', ARRAY['laptop', 'mouse']); -
Copy the customers_w directory from the customers bucket on the S3 host to the local /tmp directory and view the contents of the files:
$ java -jar /tmp/orc-tools-2.2.2-uber.jar meta -d /tmp/customers_wThe output should look as follows:
Processing data file file:/tmp/customers_w/204-0000000021_0.orc [length: 600] {"id":2,"name":"John Doe","ordered_items":["keyboard","microphone"]} ______________________________________________________________________________ Processing data file file:/tmp/customers_w/204-0000000021_1.orc [length: 587] {"id":1,"name":"Bob Brown","ordered_items":["monitor","webcam"]} ______________________________________________________________________________ Processing data file file:/tmp/customers_w/204-0000000021_3.orc [length: 602] {"id":3,"name":"Alice Green","ordered_items":["headphones","pad"]} ______________________________________________________________________________ Processing data file file:/tmp/customers_w/204-0000000021_2.orc [length: 586] {"id":4,"name":"Jane Smith","ordered_items":["laptop","mouse"]} ______________________________________________________________________________