datalake_fdw: read and write Parquet through Arrow - #1951
Conversation
372580c to
bd3f0bb
Compare
The format layer had an interface and no implementation. This is the Parquet one, and the conversions either side of it: PostgreSQL tuples into an Arrow batch, and an Arrow batch back into Datums. Arrow rather than libparquet alone, because libparquet is written in terms of Arrow's types, so linking one links the other. It is a build dependency now, found with pkg-config -- and whatever -std= its .pc file asks for is filtered out, because those flags land after CXXFLAGS and the Arrow project's own packages say -std=c++11. A fragment is a range of row groups rather than a whole file, so several segments can read one large file. Reading is single-threaded: a worker thread that fails has no way to report it through PostgreSQL. A reader and a writer hold a descriptor and memory from Arrow's allocator, which transaction abort does not reclaim, so both register with the resource owner -- the shape PAX uses in comm/pax_resource.cc. Types: bool, the four integers, both floats, text, varchar, char, bytea, date, timestamp, timestamptz. Anything else is refused by name. The 30-year epoch shift is range checked in both directions: PostgreSQL's range runs past what Arrow holds as microseconds from 1970, and a round trip through two unchecked halves would agree with itself. datalake_fdw_test is a second extension in the same library with the two functions this can be run with from SQL. Its regression case round-trips every supported type and checks that reading row groups separately gives back what reading the whole file does. Arrow 9.0.0 (EPEL 9) builds and passes; 17.0.0 (Rocky 10, gcc 14) and 17.0.0 and 21.0.0 (Rocky 8, gcc 8) compile without warnings. What is written here reads back correctly in pyarrow 21, which is not the implementation that wrote it.
bd3f0bb to
777ba63
Compare
leborchuk
left a comment
There was a problem hiding this comment.
All looks promising, but some aspects need attention or discussion.
We have working (it's not the best solution but working in production environment) extension to read iceberg for GP6 https://github.com/lithium-tech/tea
I compared data type conversions in tead and datalake_fdw. Here the differences worth attention:
┌─────┬──────────────────────┬──────────────────────────────────┬────────────────────────────────────────┬───────────────────────────────────────────┐
│ # │ Difference │ tea │ PR │ Why it deserves attention │
├─────┼──────────────────────┼──────────────────────────────────┼────────────────────────────────────────┼───────────────────────────────────────────┤
│ │ Server-encoding │ Pluggable CharsetConverter: │ │ Silently assumes server_encoding = UTF8; │
│ 1 │ conversion of │ identity, pg_custom_to_server, │ cstring_to_text_with_len on raw file │ no pg_verifymbstr in either. Produces │
│ │ strings │ or iconv UTF-8→CP1251 │ bytes (arrow_decode.c:272) │ invalid text datums. │
│ │ │ (bridge.cpp:131-179) │ │ │
├─────┼──────────────────────┼──────────────────────────────────┼────────────────────────────────────────┼───────────────────────────────────────────┤
│ 2 │ NUMERIC │ Full decimal128 → NumericVar │ Refused │ The most common column type in real lake │
│ │ │ (bridge.cpp:269-277) │ │ tables; tea also shows the typmod traps. │
├─────┼──────────────────────┼──────────────────────────────────┼────────────────────────────────────────┼───────────────────────────────────────────┤
├─────┼──────────────────────┼──────────────────────────────────┼────────────────────────────────────────┼───────────────────────────────────────────┤
│ │ │ By name (converter.h:41-53) and │ │ Positional mapping is wrong under any │
│ 5 │ Column resolution │ Iceberg field_id │ Strict positional, exact count match │ schema evolution; must not leak into the │
│ │ │ (bridge.cpp:390-443); absent │ (datalake_fdw_test.c:348-360) │ AM. │
│ │ │ column → NULL │ │ │
├─────┼──────────────────────┼──────────────────────────────────┼────────────────────────────────────────┼───────────────────────────────────────────┤
├─────┼──────────────────────┼──────────────────────────────────┼────────────────────────────────────────┼───────────────────────────────────────────┤
│ │ │ │ │ Legacy Spark/Hive timestamps arrive as │
│ 8 │ INT96 / non-µs │ Refused │ Refused │ ts(nanos); neither sets │
│ │ timestamps │ │ │ coerce_int96_timestamp_unit. Most likely │
│ │ │ │ │ first field complaint. │
├─────┼──────────────────────┼──────────────────────────────────┼────────────────────────────────────────┼───────────────────────────────────────────┤
│ │ │ │ │ Both are near-free: Iceberg time is │
│ 9 │ time, uuid │ Supported │ Refused │ µs-since-midnight = TimeADT with no epoch │
│ │ │ │ │ shift; uuid is 16 raw bytes = PG's │
│ │ │ │ │ layout. │
├─────┼──────────────────────┼──────────────────────────────────┼────────────────────────────────────────┼───────────────────────────────────────────┤
│ │ │ │ Supported, typmod re-applied on read │ A too-long value raises ERROR mid-scan on │
│ 10 │ varchar(n)/char(n) │ Not supported at all │ (arrow_decode.c:281-288) │ data you don't control. Also char(n) │
│ │ │ │ │ writes blank padding into the file. │
├─────┼──────────────────────┼──────────────────────────────────┼────────────────────────────────────────┼───────────────────────────────────────────┤
The detailed description what's wrong with datalake_fdw implementation:
#1 — encoding. This is the difference I'd raise first, because it's invisible in the PR's tests and impossible to retrofit quietly. Iceberg and Parquet define string as UTF-8. tea treats that as a conversion problem: MakePgConverter/MakeIconvConverter (bridge.cpp:447-451), InitializeIconv refusing any server encoding other than UTF-8 and WIN1251 (bridge.cpp:453-469), and iconv replacing untranslatable code points with ? rather than failing the scan. The PR copies file bytes straight into a text datum. In a WIN1251 or LATIN1 database that stores bytes no server-encoding function can interpret — upper(), length(), and text output then misbehave or throw "invalid byte sequence for encoding", possibly long after the scan. And neither side calls pg_verifymbstr, so even in a UTF8 database a malformed file injects invalid text; PG's own text input path always verifies. Minimum ask: refuse a non-UTF8 server_encoding explicitly rather than mis-decoding, and verify on the way in.
#2 — NUMERIC, and what tea got wrong doing it. Refusing DECIMAL for now is reasonable (the PR says decimal has four storage forms). What's useful is that tea's implementation has two traps the PR will hit: PGToArrowField computes precision = ((atttypmod - 4) >> 16) & 65535 (validate.cpp:59-61), which for unconstrained numeric (typmod -1) yields precision 65535 and scale 65531 — so bare numeric columns silently can't match anything. And PG allows precision up to 1000 while decimal128 caps at 38, which that expression doesn't check either. Both are one-line guards if written deliberately the first time.
#5 — positional vs field-id. Worth flagging now even though it's out of scope for a format layer, because it's the kind of decision that gets inherited. tea does two things the PR structurally can't: resolves by name against the batch schema (GetFieldIndex(parquet_name), converter.h:46) and by Iceberg field_id, and leaves a column absent from an older file as NULL rather than an error (converter.h:47 skips, isnull stays set). Both are spec requirements for add column / rename column. The PR's n_children != natts → error is right for a self-written test file and wrong for an Iceberg table.
#8 — the cheap one. Both require microseconds, so both refuse the timestamps that legacy Spark, Hive and Impala actually wrote: INT96, which Arrow's Parquet reader surfaces as timestamp(NANO). tea gets away with it because its tables are Iceberg-native. The PR is aimed at foreign files, so properties.set_coerce_int96_timestamp_unit(arrow::TimeUnit::MICRO) next to the existing property calls at parquet_read.cpp:214-222 is worth asking for — millisecond timestamps need a real decision, but INT96 is a one-liner.
| # re2-devel and parquet-devel needs thrift-devel, and on EL8 both | ||
| # of those live in EPEL. | ||
| dnf install -y \ | ||
| https://apache.jfrog.io/artifactory/arrow/almalinux/8/apache-arrow-release-latest.rpm |
There was a problem hiding this comment.
After merge we need add it to building docker container, like it was with a PAX
| * under the License. | ||
| * | ||
| * dl_resource.c | ||
| * Cleanups that happen even when nothing calls them. |
There was a problem hiding this comment.
Sorry, but here we implemented linked list. Why now reuse lib/ilist.h instead of our own implementation?
The server already ships intrusive linked lists (src/include/lib/ilist.h), and this exact pattern already uses them: contrib/pax_storage/src/cpp/comm/pax_resource.cc:35 stores a dlist_node in the entry struct and uses dlist_push_tail / dlist_delete. The PR's dl_resource.c instead hand-rolls a singly-linked list with **link splice traversal in all three functions.
The overall code will be like
typedef struct DlResourceEntry
{
dlist_node node;
ResourceOwner owner;
DlResourceRelease release;
void *arg;
} DlResourceEntry;
static dlist_head dl_resources = DLIST_STATIC_INIT(dl_resources);
/* in the callback */
dlist_foreach_modify(iter, &dl_resources)
{
DlResourceEntry *entry = dlist_container(DlResourceEntry, node, iter.cur);
if (entry->owner != CurrentResourceOwner)
continue;
if (isCommit)
elog(WARNING, "datalake_fdw leaked a resource: %p", entry->arg);
dlist_delete(&entry->node);
entry->release(entry->arg);
free(entry);
}
|
Also here are some issues, I do not know should they fixed here or in a future PR's. But I think they are quite important to be written. We could add open issues to fix them later:
|
What does this PR do?
contrib/datalake_fdwlanded in #1842 with a format layer that had noimplementation. This adds the Parquet one, and the conversions either side of
it: PostgreSQL tuples into an Arrow batch, and an Arrow batch back into Datums.
Arrow rather than libparquet alone, because libparquet is written in terms of
Arrow's types, so linking one links the other.
Decisions worth a look:
can read one large file -- the granularity the parallel scan will need.
report it through PostgreSQL's error handling.
never calls a PostgreSQL allocator; the reader is C and reads the buffers
directly, so an allocation failure unwinds through C frames only.
transaction abort does not reclaim, so both register with the resource
owner -- the shape PAX uses in
comm/pax_resource.cc.range runs past what Arrow holds as microseconds from 1970, and a round trip
through two unchecked halves would agree with itself.
Types:
bool, the four integers, both floats,text,varchar,char,bytea,date,timestamp,timestamptz. Anything else is refused by name.Type of Change
Test Plan
datalake_fdw_testis a second extension in the same library with two functions-- write a query's result to a file, read a file back as rows. Until the access
method is finished there is no other way to run this layer in a real backend.
format_parquetcategory that round-tripsevery supported type, nulls and values on both sides of 1970 included, and
checks that reading row groups 0, 1 and 2 separately gives back what
reading the whole file does. Plus the refusals: unsupported column type,
wrong type on read, column count mismatch, a range past the end of the
file, a timestamp Arrow cannot hold, a bad row group size.
make installcheck-- 4/4, on a three-segment cluster with themodule preloaded.
make -C src/test installcheck-cbdb-parallel(not run)Beyond the suite:
wrote it (9.0.0). The Arrow schema is deliberately not stored in the file, so
what comes back is what any other reader sees rather than a note of our own.
and 17.0.0 and 21.0.0 on Rocky 8 with gcc 8, compile without warnings --
compile-only checks, not test runs.
Impact
Dependencies: new build dependency on the Arrow and Parquet C++ libraries.
The module is off by default and not in the RPM, so packaging is unchanged; the
CI job that builds it with PGXS installs them -- from EPEL on Rocky 9 and 10,
and from the Arrow project's own repository pinned to 17.0.0 on Rocky 8, where
EPEL's
libarrow-develcannot be installed (itsutf8proc-develismodular-filtered out of PowerTools) and the newest Arrow wants C++20.
Arrow's
.pcfile asks for-std=c++11, and pkg-config's cflags land afterCXXFLAGS, so it is filtered out. Otherwise Arrow's headers fail to compileagainst themselves, in a way that reads like the library needing a newer
compiler.
User-facing changes: one setting,
iceberg.batch_rows. Nothing else isreachable yet -- the access method still refuses anything that would touch data.
Checklist
datalake_parquet_write(path, query)names a path on the server's file systemand runs a query through SPI, so it is as privileged as
pg_read_server_filesand granted the same way:
REVOKE EXECUTE ... FROM PUBLIC, superuser only.Additional Context
Left for later, deliberately rather than by oversight:
arrow::io::RandomAccessFileovercommon/file_system_wrapper.h, and the twoparquet files are the only ones that change when it does.
open_readerrefuses a filter set rather thanignoring it: ignoring it would still give the right rows, which is exactly why
a caller that believed pruning had happened could never find out.
change.
datalake_fdw_test's control file still lands in the share directory. PGXS'sNO_INSTALLis per-module, and these functions have to live in the librarywhose internals they test.