Pipeloom Docs
ConnectorsDestinations

S3

Set up the S3 destination connector.

Sync modes, namespaces and the columns Pipeloom adds are explained once in Connector concepts; workspace variables and custom components in Orchestration.

The S3 destination writes each stream as files in an Amazon S3 bucket, or in an S3-compatible store such as MinIO. It can write Avro, CSV, JSON Lines or Parquet files, and each stream gets a folder of its own.

This page covers setting up the bucket and access, creating the destination, and what the output looks like.

Before you start

You need:

  1. Network access from Pipeloom to your S3 storage. If the storage restricts access by IP address or sits in a private network, allow Pipeloom to connect.
  2. Encryption of data in transit enforced on the bucket.
  3. An S3 bucket, and an access key with read and write permissions on it.

These fields are always required: S3 Bucket Name, S3 Bucket Path, S3 Bucket Region and Output Format.

Step 1: Set up S3

Sign in to your AWS account, and create the bucket to write to if you don't have one yet.

Note: if your S3 storage isn't set up for TLS, connections to it silently fall back to unencrypted. We recommend TLS/SSL for every connection, in line with AWS's shared responsibility model.

Create a bucket policy

  1. Open the IAM console.
  2. Choose Policies, then Create Policy.
  3. Open the JSON tab and paste the policy below, with your bucket name in place of YOUR_BUCKET_NAME:
{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Effect": "Allow",
      "Action": [
        "s3:PutObject",
        "s3:GetObject",
        "s3:DeleteObject",
        "s3:PutObjectAcl",
        "s3:ListBucket",
        "s3:ListBucketMultipartUploads",
        "s3:AbortMultipartUpload",
        "s3:GetBucketLocation"
      ],
      "Resource": [
        "arn:aws:s3:::YOUR_BUCKET_NAME/*",
        "arn:aws:s3:::YOUR_BUCKET_NAME"
      ]
    }
  ]
}

:::note Object-level permissions alone aren't enough for the connection to work. Keep the bucket-level permissions shown above. :::

  1. Give the policy a descriptive name, then choose Create policy.

Create an IAM user and access key

Pipeloom signs in with an access key ID and secret access key. We recommend a dedicated IAM user for Pipeloom.

  1. In the IAM console, choose Users. Pick an existing user, or create one with Add users.
  2. For an existing user, open the Add permissions menu and choose Add permissions. For a new user, the permissions screen opens after you name it.
  3. Choose Attach policies directly, tick the policy you just created, then choose Next and Add permissions.
  4. Open the user's Security credentials tab and choose Create access key. Pick a use case, add tags if you want, and choose Create access key.

Keep the access key ID and secret access key for the next step.

:::note Signing in by assuming an IAM role (Role ARN) isn't available in Pipeloom yet. Leave Role ARN empty and use an access key. :::

Step 2: Create the S3 destination in Pipeloom

Create a new destination and choose S3. You can do this from the Destinations page, or from a Destination step on a pipeline canvas.

Fill in the form:

  • Access Key ID and Secret Access Key: the key from Step 1. The user needs read and write permissions on the bucket's objects.

  • S3 Bucket Name: the bucket to write to. See creating a bucket.

  • S3 Bucket Path: the folder in the bucket to write into, for example data_sync/test.

  • S3 Bucket Region: the bucket's region. See AWS region codes.

  • Output Format: Avro, CSV, JSON Lines or Parquet, with its options. See Output formats.

  • S3 Path Format (optional): how files are organized under the bucket path. The default is ${NAMESPACE}/${STREAM_NAME}/${YEAR}_${MONTH}_${DAY}_${EPOCH}_. See File paths.

  • S3 Endpoint (optional): leave empty for Amazon S3. For MinIO or another S3-compatible store, enter its URL, for example http://localhost:9000.

  • File Name Pattern (optional): how output files are named. These placeholders are supported: {date}, {date:yyyy_MM}, {timestamp}, {part_number} and {sync_id}. Don't use spaces or other placeholders; they aren't recognized.

Save the destination. Pipeloom tests the connection to the bucket, and the destination is ready once the test passes.

File paths

With the default S3 Path Format, ${NAMESPACE}/${STREAM_NAME}/${YEAR}_${MONTH}_${DAY}_${EPOCH}_, a file's full path is:

<bucket-name>/<source-namespace-if-exists>/<stream-name>/<upload-date>_<epoch>_<partition-id>.<format-extension>

For example:

testing_bucket/data_output_path/public/users/2021_01_01_1234567890_0.csv.gz
↑              ↑                ↑      ↑     ↑          ↑          ↑ ↑
|              |                |      |     |          |          | format extension
|              |                |      |     |          |          unique incremental part id
|              |                |      |     |          milliseconds since epoch
|              |                |      |     upload date in YYYY_MM_DD
|              |                |      stream name
|              |                source namespace (if it exists)
|              bucket path
bucket name

This layout means:

  1. Each stream has its own folder.
  2. Files sort by upload time.
  3. The upload time has a date part and a milliseconds part, so it's both readable and unique.

You can build your own path from these variables:

  • ${NAMESPACE}: the stream's namespace, from the source or the sync's namespace settings.
  • ${STREAM_NAME}: the stream's name.
  • ${YEAR}, ${MONTH}, ${DAY}, ${HOUR}, ${MINUTE}, ${SECOND}, ${MILLISECOND}: when the sync wrote the file.
  • ${EPOCH}: when the sync wrote the file, in milliseconds since the epoch.
  • ${UUID}: a random UUID.

Notes:

  • Repeated / characters in the path are collapsed into one.
  • The part ID is sequential, unless the bucket holds too many files, in which case it's a UUID.
  • The stream name includes the sync's stream prefix, if you set one.
  • A sync can write several files per stream: files are split by size, aiming for 200 MB compressed or less.

Sync modes

It supports namespaces: the namespace becomes part of the file path.

:::warning Full refresh overwrite deletes files In full refresh overwrite mode, a successful sync keeps the files it just wrote and deletes all earlier data for the stream. If syncs fail between runs, files from several runs may remain until the next successful sync. Each object is tagged with x-amz-meta-ab-generation-id to show which run wrote it. We recommend a bucket, or at least a path, used only by this destination, so a misconfiguration can't delete other data. :::

Output formats

Each stream is written to a folder of its own, which you can think of as the stream's table: its data is all the files in that folder.

  • In full refresh mode, old files are removed before new ones are written.
  • In incremental append mode, each sync adds new files with only the new data.

Avro

Apache Avro stores data in a compact binary format. Files always use Avro's binary encoding, and all records in a file are assumed to share one schema.

Compression codecs:

  • No compression
  • deflate, with a compression level from 0 to 9 (default 0). Level 0 is no compression and fastest; level 9 compresses best and is slowest.
  • bzip2
  • xz, with a compression level from 0 to 9 (default 6):
    • Levels 0–3 are fast, with medium compression.
    • Levels 4–6 are fairly slow, with high compression.
    • Levels 7–9 are like 6 but use bigger dictionaries and more memory. Unless a file is bigger than 8 MiB, 16 MiB or 32 MiB uncompressed, levels 7, 8 and 9 respectively just waste memory.
  • zstandard, with a compression level from -5 to 22 (default 3):
    • Negative levels are fast modes, similar to lz4 or snappy.
    • Levels above 9 are generally for archiving.
    • Levels above 18 use a lot of memory.
    • Include checksum adds a checksum to each data block.
  • snappy

Each record's JSON schema is converted to an Avro schema, and the record to an Avro record. Because the data can come from any source, that conversion has fixed rules and limitations: for example, properties missing from the stream's schema are dropped. See JSON to Avro conversion.

CSV

A CSV file has four metadata columns, plus either one column holding the whole record as JSON, or, with Root level flattening, one column per top-level field.

ColumnWhenContents
_airbyte_raw_idAlways.A UUID assigned to each record.
_airbyte_extracted_atAlways.When the record was read from the source.
_airbyte_generation_idAlways.An integer that increases with each new refresh.
_airbyte_metaAlways.A structured object with metadata about the record.
_airbyte_dataWith No flattening: the whole record is here as JSON.
root level fieldsWith Root level flattening: each top-level field gets its own column.

_airbyte_meta has these fields:

Field NameTypeDescription
changeslistA list of change objects.
sync_idintegerThe ID of the sync run.

Each change object has these fields:

Field NameTypeDescription
fieldstringThe field that was changed.
changestringThe kind of change (for example NULLED or TRUNCATED).
reasonstringWhy it changed, including where the change came from (the source, the destination, or the platform).

For example, given this record from a source:

{
  "user_id": 123,
  "name": {
    "first": "John",
    "last": "Doe"
  }
}

With no flattening, the CSV is:

_airbyte_raw_id_airbyte_extracted_at_airbyte_generation_id_airbyte_meta_airbyte_data
26d73cde-7eb1-4e1e-b7db-a4c03b4cf206162213580500011{"changes":[], "sync_id": 10111 }{ "user_id": 123, name: { "first": "John", "last": "Doe" } }

With root level flattening, the CSV is:

_airbyte_raw_id_airbyte_extracted_at_airbyte_generation_id_airbyte_metauser_idname.firstname.last
26d73cde-7eb1-4e1e-b7db-a4c03b4cf206162213580500011{"changes":[], "sync_id": 10111 }123JohnDoe

Files can be compressed with GZIP (the default); compressed files get an extra extension, .csv.gz.

JSON Lines (JSONL)

JSON Lines is a text format with one JSON object per line. Each line looks like this:

{
  "_airbyte_raw_id": "<uuid>",
  "_airbyte_extracted_at": "<timestamp>",
  "_airbyte_generation_id": "<generation-id>",
  "_airbyte_meta": "<json-meta>",
  "_airbyte_data": "<json-data-from-source>"
}

With Root level flattening, each top-level field of the record gets its own key instead of being nested under _airbyte_data.

For example, given these two records from a source:

[
  {
    "user_id": 123,
    "name": {
      "first": "John",
      "last": "Doe"
    }
  },
  {
    "user_id": 456,
    "name": {
      "first": "Jane",
      "last": "Roe"
    }
  }
]

the output file contains:

{ "_airbyte_raw_id": "26d73cde-7eb1-4e1e-b7db-a4c03b4cf206", "_airbyte_extracted_at": "1622135805000", "_airbyte_generation_id": "11", "_airbyte_meta": { "changes": [], "sync_id": 10111 }, "_airbyte_data": { "user_id": 123, "name": { "first": "John", "last": "Doe" } } }
{ "_airbyte_raw_id": "0a61de1b-9cdd-4455-a739-93572c9a5f20", "_airbyte_extracted_at": "1631948170000", "_airbyte_generation_id": "12", "_airbyte_meta": { "changes": [], "sync_id": 10112 }, "_airbyte_data": { "user_id": 456, "name": { "first": "Jane", "last": "Roe" } } }

Files can be compressed with GZIP (the default); compressed files get an extra extension, .jsonl.gz.

Parquet

Parquet output has these options:

ParameterTypeDefaultDescription
compression_codecenumUNCOMPRESSEDCompression algorithm: UNCOMPRESSED, SNAPPY, GZIP, LZO, BROTLI, LZ4 or ZSTD.
block_size_mbinteger128 (MB)Block size (row group size) in MB: how much of a row group is buffered in memory, which caps memory use while writing. Larger values make reads faster but use more memory while writing.
max_padding_size_mbinteger8 (MB)Max padding size in MB: the most padding allowed to align row groups. It's also the smallest a row group can be.
page_size_kbinteger1024 (KB)Page size in KB, used for compression. A block is made of pages, and a page is the smallest unit that has to be read in full to get one record. Too small a value makes compression worse.
dictionary_page_size_kbinteger1024 (KB)Dictionary page size in KB. With dictionary encoding there's one dictionary page per column per row group; this works like the page size, for the dictionary.
dictionary_encodingbooleantrueDictionary encoding: whether dictionary encoding is on.

These map to Parquet's ParquetOutputFormat; see its Java doc. The Parquet documentation recommends a 512–1024 MB block size and an 8 KB page size.

Each record is converted to an Avro record first, and then written as Parquet, so the same conversion rules and limitations apply.

For Parquet, the access key's user needs access to both the bucket and everything in it. Use this policy:

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Effect": "Allow",
      "Action": "s3:*",
      "Resource": [
        "arn:aws:s3:::YOUR_BUCKET_NAME/*",
        "arn:aws:s3:::YOUR_BUCKET_NAME"
      ]
    }
  ]
}

Limitations

Encryption at rest

Server-side encryption with S3-managed keys works as is. Customer-provided keys and KMS keys aren't supported.

S3-compatible stores

Not every S3-compatible store behaves the same. Known issues:

  • Linode Object Storage doesn't return ETags correctly after setting them, and the connector relies on ETags to verify that data arrived intact. It doesn't work with this destination.

On this page