Dynamic Transformers
Given that a Transformer can accept user configuration at runtime, the possibility exists for truly dynamic transformations whose logic is determined at runtime. Currently, SDL supports the following dynamic Transformers:
-
JavaScript
-
Python 3
These Transformers may or may not be present on your environment at startup (depending on the set of Transformers approved for your environment), but if approved they can be loaded easily.
How to use them
The specific usage differs slightly for each dynamic Transformer, but in general the user is responsible for providing the transformation logic, while the platform handles transport logic (i.e. consuming and producing data via the various storage media).
Dynamic JavaScript Transformer
The user is responsible for providing a function of the form:
const handler = (obj) => {
return out
}
Where obj represents the input object. Note that this Transformer only runs on JSON records passing through Kafka.
Where input represents an incoming JSON record, and output will be unmarshalled as JSON and sent to Kafka.
Dynamic Python Transformer
This Transformer is slightly more complex; the user is responsible for providing a runnable Python 3 file (with main()). The df-daft-py client library is embedded to simplify setting up the transport logic. See the df-daft-py documentation for more information on how to use this library. Below is an example of a transformation script:
# Transforms messages from one Kafka topic to another
import os
from df_daft_py.kafka.kafka import start_transform_loop
def transform(data) -> dict:
# Make transformations here, for example:
# data["latitude"] = 0.12345
return data
def main():
src_topic = os.getenv("SOURCE_KAFKA")
dest_topic = os.getenv("DEST_KAFKA")
start_transform_loop(src_topic, dest_topic, transform)
Additionally, users may specify a requirements.txt file configuration, which will be installed before the script is run. Note that on airgapped environments, this functionality may not work correctly unless a suitable mirror repository is available.
Test Python Code Before Starting a Pipeline
When the deployed Dynamic Transformer (Python) template advertises Test Code support, the Code Editor can run the current unsaved script against sample Kafka JSON without starting the Pipeline.
The first Python Test Code profile has these rules:
-
The
scriptfile must define a callabletransform(data)function. -
Each
SOURCE_KAFKAmessage value must be a JSON object. The decoded object is passed totransform. -
transformmust return a dictionary, which is written toDEST_KAFKA, orNoneto skip that message. -
Message keys and string headers are preserved.
-
The Python standard library and packages already included in the runtime image are available.
-
requirements.txtmust be blank or contain comments only. Test Code does not download or install packages. -
MinIO input is not supported by Test Code. This does not prevent the normal live transformer from using its supported MinIO workflow.
Keep transform independent from the transport setup, then pass that same function to the live loop:
def transform(data) -> dict:
return {**data, "validated": True}
def main():
start_transform_loop(
os.getenv("SOURCE_KAFKA"),
os.getenv("DEST_KAFKA"),
transform,
)
Test Code does not call main(), connect to Kafka or MinIO, or save editor changes. It runs in a temporary isolated workload with time, memory, CPU, request, response, record, and log limits.
Troubleshooting:
-
If Test Code is unavailable, the environment has not enabled a compatible Python simulation image. Live Python execution may still be available.
-
If
requirements.txtis rejected, remove non-comment entries for the test or use a package already included in the image. -
If a MinIO input is connected, use a Kafka JSON sample for Test Code or validate the MinIO path through a live Pipeline.
-
If the result reports a missing
transform, syntax, import, runtime, or serialization error, correct the script and run the test again.