Skip to content

datalake_fdw: storage I/O over S3 and pluggable storage backends - #2044

Open
MisterRaindrop wants to merge 7 commits into
apache:mainfrom
MisterRaindrop:feature/datalake-s3
Open

MisterRaindrop wants to merge 7 commits into
apache:mainfrom
MisterRaindrop:feature/datalake-s3

Conversation

@MisterRaindrop

@MisterRaindrop MisterRaindrop commented Sep 23, 2026 •

Copy link
Copy Markdown
Contributor

What does this PR do?

contrib/datalake_fdw reads and writes Parquet since #1951, but only on a
local file system. This adds the storage layer underneath: an s3 backend
over the AWS SDK for C++, and a published contract so a backend this extension
does not ship can be added from outside it.

Closes #2009.

Decisions worth a look:

  • A backend answers one question: given a location and its options, which
    arrow::fs::FileSystem reads and writes it. Opening files, listing,
    classifying errors, accounting for memory and translating a URI into a native
    path all stay on this extension's side of the boundary, so a backend is a
    mount function and nothing else.
  • s3 is our own arrow::fs::FileSystem over the AWS SDK, not Arrow's.
    No RPM of Arrow is built with S3 support -- EPEL's and the Arrow project's
    own both set use_s3 0 -- so depending on it would mean asking every user to
    build Arrow. The SDK is a BUILD_ONLY="s3;sts" static build that takes well
    under a minute, links statically, and adds libcurl, OpenSSL and zlib to
    NEEDED and nothing else.
  • Timeouts and retries are set by this module, not inherited: connect 5 s,
    request 300 s, three retries, so a black-holed endpoint answers in about 24
    seconds rather than hanging a session, on every Arrow version.
  • A writer only ever removes what it created. The local backend opens with
    O_EXCL and unlinks that path alone; the S3 one keys cleanup on whether an
    upload is still live rather than on whether the stream is closed, so a
    Close() that fails at the last part still aborts the upload instead of
    leaving parts to be billed for.
  • Credentials are scrubbed at one place. Every value under a key containing
    secret, token or password is remembered per process, and dl_error_set
    -- the single exit every DlErrCode error goes through -- removes it from
    the message. An unexpected C++ exception and a third-party backend's own
    status are covered by construction rather than by remembering to call
    something. A base_path that carries a password in its userinfo is redacted
    the same way, in the message as well as the detail. Where one remembered
    secret is a prefix of another, the longer is the one masked; and a ? or
    # inside the userinfo of a malformed URI does not stop the redactor short
    of the @ that hides the password.
  • Which schemes a volume may name comes from the registry, not a list in
    the parser: a backend that registered mine can be written into a
    base_path without this extension being changed.
  • The plug-in contract is versioned three ways: an ABI version, an Arrow
    version, and a fingerprint over the compiler major and _GLIBCXX_USE_CXX11_ABI.
    Registration happens during preload and in any order -- the registering side
    pulls in datalake_fdw through load_external_function, inside a PG_TRY
    so an ereport cannot unwind through the plug-in's C++ frames.
  • Every Arrow allocation still goes through the tracked pool, including a
    backend's: mount is handed a host struct carrying it, and the five
    installed headers give a backend no way to reach Arrow's default pool by
    accident.

Six commits, each buildable and green on its own: the storage layer, the s3
backend, the Parquet reader and writer moving onto the facade, the conformance
suite, the README, and the CI wiring.

Type of Change

  • Bug fix (non-breaking change)
  • New feature (non-breaking change)
  • Breaking change (fix or feature with breaking changes)
  • Documentation update

Test Plan

One body of storage behaviour runs against every backend rather than once per
backend: storage_conformance parameterises a shared SQL file over a URI
prefix and a volume, and runs it for file://, for s3:// and for the test
backend.

  • Integration tests added -- three new categories. storage_local covers
    the facade and the registry: four kinds of bad registration, a duplicate
    scheme, path escapes, a name already in use, a missing file.
    storage_s3 covers the backend against a real service: round trip,
    9 MB through multipart, ListObjectsV2 paging, HeadObject,
    DeleteObject, deleting twice, a missing key, a missing bucket, a wrong
    secret (and that the message does not contain it), path style with an
    endpoint override, and a black-holed endpoint inside 30 s.
    storage_conformance adds the type round trip, field-id projection, a
    row-group range, listing, a rejected overwrite, a write that fails
    halfway leaving nothing behind, a prefix nothing was written under and
    one whose objects were just deleted (both not-found, which is the point:
    object storage and a filesystem disagree and the facade decides), and
    volume resolution: no USAGE, a PUBLIC mapping, a URI outside its
    volume, a volume that does not exist, and a rejected base_path that
    must not echo its own password.
  • Passed make installcheck -- every category, on a three-segment cluster
    with the module preloaded, against MinIO and against SeaweedFS, on
    Arrow 9.0.0 and Arrow 17.0.0. Each of the six commits was built and
    checked separately, so the series bisects. The server that runs the
    suite has to preload the test backend as well -- see Running the
    tests
    below.
  • Unit tests added/updated
  • Passed make -C src/test installcheck-cbdb-parallel (not run)

Beyond the suite:

  • The plug-in contract was used from outside. A .cc file that includes
    only the five installed headers and PostgreSQL's server headers compiles to a
    .so that registers a backend; put before datalake_fdw in
    shared_preload_libraries the cluster starts and the scheme works, and
    built with a stale abi_fingerprint the registration is refused with both
    values in the message.
  • Memory was measured, S3 against local. Writing and reading the same
    Parquet file through s3:// costs about 15 MiB of resident memory more than
    through file:// -- and that difference stays flat as the object grows
    (7.9 MiB for a 2 MB file, 14.8 MiB for a 149 MB one), so neither the response
    body nor the multipart buffer accumulates.
  • The credential chain was exercised. With no user mapping the same
    statement fails; with AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY in the
    postmaster's environment it succeeds.
  • A multipart upload that fails is aborted, and that is checked outside
    the suite.
    No SQL can see an abandoned upload -- listing a prefix does
    not show one -- so a mutation that deletes the AbortMultipartUpload call
    still passes installcheck. What does catch it is the server's own list of
    in-flight uploads: one entry with the mutation, none without it, on MinIO
    and on SeaweedFS alike. CI therefore asks SeaweedFS for its .uploads after
    the run and fails if anything is there.
  • A ranged read survives a transient error. With a proxy in front of
    MinIO answering the first one to four ranged GETs with 503 SlowDown, the
    read recovers after one, two and three of them with the bytes intact, and
    after four -- the request and its three retries -- fails with an error.
    Each attempt writes the body from the start of the destination, so an
    error body left by a failed attempt cannot crowd out the retry.
  • A listing can be cancelled. Against a fake endpoint that answers every
    ListObjectsV2 page with "more", pg_cancel_backend ends the statement in
    about 0.1 s with the usual canceling statement due to user request; the
    page loop checks for a pending interrupt and hands it back to C, since it
    cannot be serviced from a C++ frame.
  • It builds --with-llvm. With bitcode generation on, the module and the
    test backend both build, every .bc included.
  • Nothing new is exported. nm -D lists the PostgreSQL entry points, one
    registration function and the test extension's UDFs; no Arrow and no AWS
    symbols, on all three build variants.

Impact

Dependencies: an optional build dependency on the AWS SDK for C++. Without
it the extension builds as before and says so, and opening an s3:// location
reports that it was built without it; naming a prefix that has no SDK in it is
an error rather than a silent fallback. The module is off by default and not in
the RPM, so packaging is unchanged.

CI builds the SDK from source once per distribution and architecture and caches
it, starts SeaweedFS for the s3 cases, and now takes Arrow from the Arrow
project's own repository on Rocky 9 and 10 as well, pinned -- 17.0.0 and
21.0.0, so a version change in EPEL cannot arrive without a commit.

What that costs the ic-datalake-fdw job, measured on all three legs: building
the SDK takes 175-228 s on a cache miss and the cache restores in 1 s (it is
4 MB); SeaweedFS is up 6 s after the step starts; the s3 cases add about 35 s,
of which 22 s is one deliberate connect timeout.

Running the tests: the test backend is a library of its own,
datalake_fdw_dltest, built and installed beside datalake_fdw but never
loaded by it -- so production clusters have no test scheme. A server that runs
installcheck needs both:

shared_preload_libraries = 'datalake_fdw,datalake_fdw_dltest'

make check and CI already set this. Without it storage_local and
conformance_dltest fail, because the scheme they test is not registered.

User-facing changes: a volume's base_path now accepts file:// as well
as s3://, and the scheme it rejects is named against what is registered
rather than against a fixed list. A . or .. path segment is rejected
anywhere in a location, so a URI cannot climb out of the volume it names.
Three iceberg_am_reject DETAIL lines and
two ERROR lines change wording, the latter because they used to quote back a
URI that can carry a password.

Checklist

contrib/datalake_fdw/README.md documents volumes, where credentials come
from, building with S3 support, and writing a backend -- the example in it is
compiled against the installed headers rather than written out by hand. The
test extension's new functions name a path on the server's file system and are
revoked from PUBLIC like the two that were already there.

Additional Context

Known limits, deliberately rather than by oversight:

  • A file:// volume must be the same directory on every host. Nothing
    checks it; the README says so.
  • Between checking that a name is free and creating it, another writer can
    take it.
    Iceberg's file names are unique by construction, so it does not
    arise there.
  • An upload abandoned by a crashed backend leaves its parts behind. A
    bucket lifecycle rule that expires incomplete multipart uploads is the usual
    answer; the README says so.
  • About 15 MiB per backend process is outside gp_vmem_protect_limit --
    the SDK client, its connection and one part buffer are allocated by the SDK
    rather than through the tracked pool.
  • The plug-in contract is a source-level interface, not a binary-stable
    ABI.
    mount returns a struct holding a std::shared_ptr and a
    std::string; the fingerprint catches the mismatches that occur in practice
    (compiler major, the libstdc++ dual ABI, Arrow version), not every possible
    one. A backend is supported when it is rebuilt with the same compiler major
    and Arrow package as this module.
  • One S3 object tops out at about 78 GiB. Parts are a fixed 8 MiB and the
    10,000-part limit is left to the server to enforce.
  • CI does not verify digests of the SeaweedFS binary and the AWS SDK
    sources it downloads; both are pinned by version only.
  • No HDFS. The scheme is gone from the parser rather than half-supported;
    it returns as a backend when someone needs it.

@MisterRaindrop
MisterRaindrop force-pushed the feature/datalake-s3 branch 3 times, most recently from dbd0f51 to eb15606 Compare September 24, 2026 10:26
Where a lake table's files live was decided by a switch on the URI
scheme, and only s3 and hdfs were written into it.  This replaces that
with a contract: a backend answers one question -- given a location and
its options, which arrow::fs::FileSystem reads and writes it -- and
everything else stays on this extension's side of the boundary.

What lands here:

- The contract, in five headers that "make install" puts under
  $(includedir_server)/extension/datalake_fdw/, so a backend can be
  built outside this tree against nothing else.  It carries an ABI
  version, the Arrow version, and a fingerprint over the compiler major
  and _GLIBCXX_USE_CXX11_ABI, because shared_ptr and arrow::Result
  cross the boundary by value.
- A registry keyed by URI scheme, found through a rendezvous variable.
  Registration happens during preload and in any order: the registering
  side pulls in datalake_fdw through load_external_function, inside a
  PG_TRY so an ereport cannot unwind through a plug-in's C++ frames.
  Each process initializes a backend at its first mount and registers
  the matching finalizer then, because on_exit_reset() clears what a
  postmaster child inherits.
- The facade over that filesystem: opening, listing, deleting, turning
  a URI into a native path, and classifying failures into DlErrCode.
  Backends never re-parse a URI, and an empty listing is ruled
  not-found here rather than by each backend, because object storage
  and a filesystem disagree about what an empty prefix means.  A
  backend that classified a failure itself is believed; anything else
  is read as a whole status rather than as a message, since Arrow puts
  a failed open's errno in a detail -- and only phrases a filesystem
  writes about itself are matched, because a service's numeric code
  found inside the text would turn an error about a path into an error
  about the path not existing.
- A file backend for shared mounts, which creates with O_EXCL and, if
  it has to give up, removes the file this writer created and nothing
  else.  The stream takes ownership before the guard is released and
  removes the file from its destructor, so an exception leaves nothing
  behind either.
- Credential scrubbing in one place.  Every value under a key
  containing secret, token or password is remembered per process, and
  dl_error_set -- the single exit every DlErrCode error goes through --
  removes it from the message, so an unexpected C++ exception and a
  third-party backend's own status are covered by construction.  No
  value is declined for being short or for being the thirty-third, and
  a message is scrubbed as it is copied rather than after: a field that
  truncates would otherwise leave the first half of a key behind.  The
  list is held with malloc, so recording a credential cannot ereport
  out of the C++ frame that is holding one; if it fails anyway, the
  message is withheld instead of shown.
- An example backend, built as a library of its own under dltest/.  It
  is not part of this one: it compiles against the installed headers,
  registers through datalake_storage_register, and a server that does
  not name it has no dltest scheme.  Building it in would put a
  test-only scheme on every cluster -- one a volume could be created
  against -- and would demonstrate nothing about the boundary it exists
  to demonstrate.

The parser now accepts file:// as well as s3://, and names what it
rejects against the registry rather than a fixed list, so a scheme a
plug-in registered can be written into a base_path.  What it quotes
back is redacted first: a rejected URI's userinfo and its query are
reported as present rather than reproduced, one being able to hold a
password and the other a presigned signature.  Userinfo is cut at the
last "@" in the authority rather than the first, because it may hold
one of its own.  Three DETAIL lines in iceberg_am_reject change with
it.

"." and ".." are refused as whole path segments, whatever the scheme.
Which paths a volume covers is decided by comparing prefixes as text,
and "file:///v/inside/../secret" satisfies that comparison while
resolving somewhere else, so the URI never reaches it.

The old s3 stub is removed here rather than adapted twice; the backend
comes back, written against this contract, in the commit that follows.
s3 returns as an arrow::fs::FileSystem of our own over the AWS SDK for
C++, rather than Arrow's S3FileSystem: no RPM of Arrow is built with S3
support -- EPEL's and the Arrow project's own both set use_s3 0 -- so
depending on it would mean asking every user to build Arrow.

A synchronous S3Client, so a cancelled query cannot leave a callback
holding a backend's memory.  Reads go through PreallocatedStreamBuf
straight into a tracked buffer rather than the SDK's default
stringstream.  The range is clamped against the size the file was
opened at, so every byte asked for existed when the request was made
and a short answer is a truncated response or an object replaced
underneath the reader -- an error, not an end of file.  Writes buffer 8 MiB and begin a multipart upload
only when they exceed it, so a small file is one PutObject.  Cleanup
depends on whether an upload is still live rather than on whether the
stream was closed, so a Close() that fails at its last part still
aborts instead of leaving parts to be billed for, and the destructor
does the same.

A listing follows continuation tokens and gives up if one stops
advancing, rather than trusting the service to end it.  Credentials are
taken as a pair or not at all -- half of one is a configuration error,
not a reason to fall back to the host's own identity -- and a bucket
that does not exist is reported as missing rather than as a directory.
The options are spelled as the DDL spells them, access_key_id and
secret_access_key and session_token, rather than under a second
vocabulary of the backend's own.

Connect timeout 5 s, request timeout 300 s, three retries, set here
rather than inherited, so a black-holed endpoint answers in about 24
seconds on every Arrow version.

The SDK is found by prefix -- the one AWS_SDK_PREFIX names, or the
usual places a hand-built one lands.  Without it the module still
builds and says so, and opening an s3:// location reports that it was
left out; a prefix with no SDK in it is an error rather than a silent
fallback.  It is linked statically even where that prefix also holds
shared copies of it, and adds libcurl, OpenSSL and zlib to NEEDED and
nothing else.
The Parquet reader and writer opened paths through the local file
system directly.  They now open through the storage facade, so a
fragment names a filesystem and a path relative to its mount, and the
same reader works against a volume on s3 as on a shared mount.

"Create only" moves with them: the writer asks the facade, which is the
one place that checks a name is free, and giving up is the stream's
Abort rather than a delete by path -- after a failed write the name may
already belong to somebody else.

A volume resolves to a location plus credentials in one place, checking
USAGE on the server and reading the user mapping, the PUBLIC mapping,
or neither.  What its validator quotes back when it rejects a base_path
goes through the same redaction as everything else that echoes one, so
a password in the userinfo does not reach the server log.
Three backends that must behave identically are worth one body of
tests, not three that drift.  A shared SQL file is parameterised over a
URI prefix and a volume and run for file://, for s3:// and for the test
backend: the type round trip, field-id projection, a row-group range,
listing, a rejected overwrite, a write that fails halfway leaving
nothing behind, a prefix nothing was ever written under, and -- where
there is a filesystem underneath -- one holding nothing but an empty
directory.  Both of those have to read as not found rather than as an
empty listing, which is the disagreement between a bucket and a
directory that the facade exists to settle.

The write that fails halfway carries thirty-two distinct md5s a row and
asks for no compression, because a kilobyte of one repeated character
encodes to nothing: twenty thousand such rows came to 174 KB, so the
case that was meant to fail in the middle of a multipart upload had
never begun one.  It now writes 14 MiB before it fails.  What the
listing afterwards proves is that no object is left; that no upload is
left running is asserted outside the suite, against the service's own
list, because an abandoned upload leaves no object either.

Volume resolution gets its own case, on a file volume so it runs
everywhere: no USAGE, a PUBLIC mapping, a URI outside its volume, a
volume that does not exist, a base_path whose userinfo or query must
not come back in the error -- including one whose userinfo holds an "@"
of its own -- and a URI that would leave the volume through "..".

Listing past one page is a few thousand requests and minutes of wall
clock, so it is a case of its own and runs where
DATALAKE_TEST_S3_PAGINATION asks for it rather than on every build.
Volumes and their options, where credentials come from and in what
order, what a file:// volume requires of the filesystem, how to build
with S3 support, and how to write a backend -- the worked example is
dltest/, which is compiled and run rather than written out by hand, and
the sketch in the text says plainly which of the obligations below it
does not meet.  Plus
the limits that are accepted rather than fixed, among them what s3
costs: a backend process holds about 15 MiB of resident memory that
gp_vmem_protect_limit does not see, and that stays flat as the object
grows -- 7.9 MiB for a 2 MB file, 14.8 MiB for a 149 MB one -- because
it is the client, its connection and one part buffer rather than a cost
per byte.
The s3 half of contrib/datalake_fdw's regression had nothing to run
against.  Three things were missing.

The AWS SDK for C++, which no distribution packages: built from source
with BUILD_ONLY="s3;sts", and cached under the version, the
distribution and the architecture, since nothing else changes it.  A
miss costs about three minutes; the cache is 4 MB and restores in one
second.  Only the build dependencies the image lacks are installed,
asked for by capability -- naming a package it already has makes dnf
try to upgrade it, and on Rocky 10 the newest libcurl-devel wants a
libcurl no enabled repository carries.

A service that speaks S3: SeaweedFS, one static binary, pinned, and
started with an identity file -- without one it accepts any
credentials, and the case that asserts a wrong secret is refused would
pass by not being tested.  Readiness waits on the S3 port itself, which
begins listening seconds after the master elects itself.

Its coordinates, which "su - gpadmin" drops along with the rest of the
environment, so they are named on the command line.  Every other test
entry leaves that empty, and a leg where the service did not start
skips the s3 cases rather than failing them.

Arrow on Rocky 9 and 10 now comes from the Arrow project's repository
pinned to 17.0.0 and 21.0.0, as it already did on Rocky 8.  EPEL's
moves when EPEL does, and Arrow is the library this extension's ABI is
shared with.

A green job is only evidence about the cases that ran, so afterwards the
three s3 case names are looked for in the log the run wrote; a variable
that failed to arrive makes the Makefile skip them and report success.
The wrong-secret value the suite asserts against is then searched for in
the segment logs and the service's own log, which is where a credential
would land that never reached the client, and the service is asked
whether any multipart upload is still in progress -- an upload nobody
aborted leaves no object to notice it by.
@MisterRaindrop
MisterRaindrop marked this pull request as ready for review October 4, 2026 01:39
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

datalake_fdw: storage I/O over object storage (S3)

1 participant