For the complete documentation index, see llms.txt. This page is also available as Markdown.

Ingestion Transformations

Raw source data often needs to undergo some transformations before it is pushed to Pinot.

Transformations include extracting records from nested objects, applying simple transform functions on certain columns, filtering out unwanted columns, as well as more advanced operations like joining between datasets.

A preprocessing job is usually needed to perform these operations. In streaming data sources, you might write a Samza job and create an intermediate topic to store the transformed data.

For simple transformations, this can result in inconsistencies in the batch/stream data source and increase maintenance and operator overhead.

To make things easier, Pinot supports transformations that can be applied via the table config.

Transformation functions

Pinot supports the following functions:

  • Groovy functions

  • Built-in functions

Groovy functions

Groovy functions can be defined using the syntax:

Groovy({groovy script}, argument1, argument2...argumentN)

Any valid Groovy expression can be used.

⚠️ Enabling Groovy

Allowing executable Groovy in ingestion transformation can be a security vulnerability. To enable Groovy for ingestion, set the following controller configuration:

controller.disable.ingestion.groovy=false

If not set, Groovy for ingestion transformation is disabled by default.

Groovy static analysis

Pinot can also apply static analysis to Groovy scripts before compiling them. This is configured with cluster-level Groovy analyzer configs:

  • pinot.groovy.all.static.analyzer: default analyzer for both query-time and ingestion-time Groovy

  • pinot.groovy.ingestion.static.analyzer: ingestion-specific override used when Pinot validates Groovy in table configs

  • pinot.groovy.query.static.analyzer: query-specific override used when Pinot validates Groovy in broker queries

If the ingestion-specific or query-specific config is not set, Pinot falls back to pinot.groovy.all.static.analyzer. If none of these cluster configs are set, Groovy static analysis is disabled.

Each static analyzer config is a JSON object with these fields:

  • allowedReceivers: fully qualified receiver classes Groovy method calls can target

  • allowedImports: fully qualified imports allowed in the script

  • allowedStaticImports: fully qualified static imports allowed in the script

  • disallowedMethodNames: method names Pinot rejects even if the receiver is otherwise allowed

  • methodDefinitionAllowed: whether Groovy method definitions are allowed inside the script

Static analysis does not enable Groovy by itself. You still need controller.disable.ingestion.groovy=false to use Groovy in ingestion transforms.

Use the controller API to inspect and update these configs:

  • GET /cluster/configs/groovy/staticAnalyzerConfig/default returns Pinot's built-in sample config

  • GET /cluster/configs/groovy/staticAnalyzerConfig returns the cluster's current Groovy static analyzer config overrides

  • POST /cluster/configs/groovy/staticAnalyzerConfig updates one or more of the Groovy analyzer configs

The full request and response examples for these endpoints are in the controller API reference.

Built-in Pinot functions

All the functions defined in this directory annotated with @ScalarFunction (for example, toEpochSeconds) are supported ingestion transformation functions.

Below are some commonly used built-in Pinot functions for ingestion transformations.

DateTime functions

These functions enable time transformations.

toEpochXXX

Converts from epoch milliseconds to a higher granularity.

Function name
Description

toEpochSeconds

Converts epoch millis to epoch seconds. Usage:"toEpochSeconds(millis)"

toEpochMinutes

Converts epoch millis to epoch minutes Usage: "toEpochMinutes(millis)"

toEpochHours

Converts epoch millis to epoch hours Usage: "toEpochHours(millis)"

toEpochDays

Converts epoch millis to epoch days Usage: "toEpochDays(millis)"

toEpochXXXRounded

Converts from epoch milliseconds to another granularity, rounding to the nearest rounding bucket. For example, 1588469352000 (2020-05-01 42:29:12) is 26474489 minutesSinceEpoch. `toEpochMinutesRounded(1588469352000) = 26474480 (2020-05-01 42:20:00)

Function Name
Description

toEpochSecondsRounded

Converts epoch millis to epoch seconds, rounding to nearest rounding bucket"toEpochSecondsRounded(millis, 30)"

toEpochMinutesRounded

Converts epoch millis to epoch seconds, rounding to nearest rounding bucket"toEpochMinutesRounded(millis, 10)"

toEpochHoursRounded

Converts epoch millis to epoch seconds, rounding to nearest rounding bucket"toEpochHoursRounded(millis, 6)"

toEpochDaysRounded

Converts epoch millis to epoch seconds, rounding to nearest rounding bucket"toEpochDaysRounded(millis, 7)"

fromEpochXXX

Converts from an epoch granularity to milliseconds.

Function Name
Description

fromEpochSeconds

Converts from epoch seconds to milliseconds "fromEpochSeconds(secondsSinceEpoch)"

fromEpochMinutes

Converts from epoch minutes to milliseconds "fromEpochMinutes(minutesSinceEpoch)"

fromEpochHours

Converts from epoch hours to milliseconds "fromEpochHours(hoursSinceEpoch)"

fromEpochDays

Converts from epoch days to milliseconds "fromEpochDays(daysSinceEpoch)"

Simple date format

Converts simple date format strings to milliseconds and vice versa, per the provided pattern string.

Function name
Description

Converts from milliseconds to a formatted date time string, as per the provided pattern "toDateTime(millis, 'yyyy-MM-dd')"

Converts a formatted date time string to milliseconds, as per the provided pattern "fromDateTime(dateTimeStr, 'EEE MMM dd HH:mm:ss ZZZ yyyy')"

Note

Letters that are not part of Simple Date Time legend (https://docs.oracle.com/javase/8/docs/api/java/text/SimpleDateFormat.html) need to be escaped. For example:

"transformFunction": "fromDateTime(dateTimeStr, 'yyyy-MM-dd''T''HH:mm:ss')"

JSON functions

Function name
Description

json_format

Converts a JSON/AVRO complex object to a string. This json map can then be queried using jsonExtractScalar function. "json_format(jsonMapField)"

Types of transformation

Filtering

Records can be filtered as they are ingested. A filter function can be specified in the filterConfigs in the ingestionConfigs of the table config.

If the expression evaluates to true, the record will be filtered out. The expressions can use any of the transform functions described in the previous section.

Consider a table that has a column timestamp. If you want to filter out records that are older than timestamp 1589007600000, you could apply the following function:

Consider a table that has a string column campaign and a multi-value column double column prices. If you want to filter out records where campaign = 'X' or 'Y' and sum of all elements in prices is less than 100, you could apply the following function:

Filter config also supports SQL-like expression of built-in scalar functions for filtering records (starting v 0.11.0+). Example:

Column transformation

Transform functions can be defined on columns in the ingestion config of the table config.

Pinot evaluates ingestion transforms against normalized extractor values, not the raw source-format encoding. For typed input formats such as Avro, Parquet, ORC, Thrift, and Protocol Buffers, booleans stay Boolean, Byte and Short values widen to Integer, logical dates and times surface as java.time.LocalDate, java.time.LocalTime, or java.sql.Timestamp, multi-value fields surface as Object[], and maps or nested records surface as Map<Object, Object>. Pinot coerces those intermediate values to the column's declared schema type after the transform step.

If you have an older transform that depended on format-specific quirks such as stringified booleans or raw epoch date and time values, update the transform to handle the normalized value or cast it explicitly.

For example, imagine that our source data contains the prices and timestamp fields. We want to extract the maximum price and store that in the maxPrices field and convert the timestamp into the number of hours since the epoch and store it in the hoursSinceEpoch field. You can do this by applying the following transformation:

Below are some examples of commonly used functions.

String concatenation

Concat firstName and lastName to get fullName

Find an element in an array

Find max value in array bids

Time transformation

Convert timestamp from MILLISECONDS to HOURS

Column name change

Change name of the column from user_id to userId

Rename fields from a Kafka JSON message

Kafka JSON payloads often use keys that aren’t great Pinot column names. Common examples are keys containing -, such as event-id.

Map the source key to a schema-friendly column using transformConfigs. Reference the source key with a quoted identifier.

Add the destination columns (for example, event_id) to your Pinot schema.

Extract value from a column containing space

Pinot doesn't support columns that have spaces, so if a source data column has a space, we'll need to store that value in a column with a supported name. To extract the value from first Name into the column firstName, run the following:

Ternary operation

If eventType is IMPRESSION set impression to 1. Similar for CLICK.

AVRO Map

Store an AVRO Map in Pinot as two multi-value columns. Sort the keys, to maintain the mapping. 1) The keys of the map as map_keys 2) The values of the map as map_values

Chaining transformations

Transformations can be chained. This means that you can use a field created by a transformation in another transformation function.

For example, we might have the following JSON document in the data field of our source data:

We can apply one transformation to extract the userId and then another one to pull out the numerical part of the identifier:

Pinot validates these chains by dependency, not by list order. A transform can consume a column produced by another transform even if the consumer appears earlier in transformConfigs, as long as the dependency graph is valid.

Intermediate transform outputs do not need schema columns when they are only consumed by later transforms. Final outputs that Pinot stores or aggregates still need schema columns.

For example, the following pattern is valid even when message_obj is not in the schema:

In this example, message_obj is an intermediate column, so Pinot materializes it during transformation and drops it before indexing. The leaf columns such as level and msg_time are the stored outputs, so those columns still belong in the schema.

Flattening

There are 2 kinds of flattening:

One record into many

This is not natively supported as of yet. You can write a custom Decoder/RecordReader if you want to use this. Once the Decoder generates the multiple GenericRows from the provided input record, a List<GenericRow> should be set into the destination GenericRow, with the key $MULTIPLE_RECORDS_KEY$. The segment generation drivers will treat this as a special case and handle the multiple records case.

Extract attributes from complex objects

The multi-stage query engine (MSE) supports flattening array columns at query time using CROSS JOIN UNNEST(...). This is the recommended approach for expanding array fields without pre-flattening data at ingestion.

For full syntax including WITH ORDINALITY, filtering, and aggregation after unnesting, see the Unnest operator page.

Add a new column during ingestion

If a new column is added to table or schema configuration during ingestion, incorrect data may appear in the consuming segment(s).

To ensure accurate values are reloaded, do the following:

  1. Pause consumption (and wait for pause status success): $ curl -X POST {controllerHost}/tables/{tableName}/pauseConsumption

  2. Apply new table or schema configurations.

  3. Resume consumption: $ curl -X POST {controllerHost}/tables/{tableName}/resumeConsumption

Last updated

Was this helpful?