File Source Connector
stream sourcescan table source
The File Source connector reads file content into the rekuiper stream processing pipeline.
The connector supports batch file processing and real-time directory monitoring. When monitoring a directory, all files in the directory must share the same format. The engine processes files in alphabetical order by filename.
Supported File Types
The file connector supports the following file structures:
- JSON: Files containing standard JSON arrays.
- CSV: Comma-separated or custom-delimited tabular text files.
- Lines: Text files with one record per line.
- Parquet: Columnar Apache Parquet files using Arrow schema inference.
- Raw: Reads the entire file as a single binary payload.
JSON Array Example
[
{"id": 1, "name": "John Doe"},
{"id": 2, "name": "Jane Smith"}
]NOTE
If a file contains multiple JSON records separated by newlines, set fileType to "lines" and FORMAT to "json".
CSV Example
id,name,age
1,John Doe,30
2,Jane Smith,25Custom separators (such as spaces or semicolons) are supported:
id name age
1 John Doe 30
2 Jane Smith 25Lines Example
Each line represents a distinct event record:
{"id": 1, "name": "John Doe"}
{"id": 2, "name": "Jane Smith"}You can combine lines with binary formats. For example, set FORMAT to "protobuf" and provide a schema to parse newline-separated Protobuf messages.
Configuration Parameters
Configure the connector in etc/sources/file.yaml:
default:
fileType: json
path: data
interval: 0
sendInterval: 0
actionAfterRead: 0
moveTo: /tmp/kuiper/moved
hasHeader: false
# columns: [id, name]
ignoreStartLines: 0
ignoreEndLines: 0
decompression: ""File Type and Directory
fileType: File format:"raw","json","csv", or"lines". When using"raw", set the stream format to"binary".path: Directory path relative to the rekuiper root directory, or an absolute path. Do not include filenames here; specify filenames inDATASOURCE.
Reading and Sending Intervals
interval: Polling interval in milliseconds. Settinginterval: 0activates filesystem change monitoring instead of polling. When files change or new files appear, the engine reads them immediately.sendInterval: Delay in milliseconds between emitted event records.
Post-Read Actions
actionAfterRead: Defines file handling after reading completes:0: Keep the file.1: Delete the file.2: Move the file to the path specified inmoveTo.
moveTo: Target directory path whenactionAfterReadis set to2.
CSV Parsing Options
hasHeader: Boolean indicating whether the first row contains column headers.columns: List of column names when files lack headers (for example,columns: [id, name]).ignoreStartLines: Number of lines to skip at the beginning of the file. Empty lines are ignored and not counted.ignoreEndLines: Number of lines to skip at the end of the file.
Decompression
decompression: Decompresses incoming archive files. Supported algorithms:"gzip"and"zstd".
Create a Table Source
The file source commonly operates as a scan table for static reference lookups:
CREATE TABLE table1 (
name STRING,
size BIGINT,
id BIGINT
) WITH (DATASOURCE = "lookup.json", FORMAT = "json", TYPE = "file");Create a rule that joins stream events with the file table:
CREATE RULE rule1 AS SELECT * FROM fileDemo WHERE temperature > 50 INTO mySink;To manage tables through the REST API or CLI, refer to Tables Management with REST API and Tables Management with CLI.
Configuration Tutorials
Tutorial 1: Parse Space-Delimited CSV Files
Consider this space-separated data file:
id name age
1 John 56
2 Jane 34Configure
etc/sources/file.yaml:yamlcsv: fileType: csv hasHeader: trueCreate a stream using the
DELIMITEDformat and specify a space character delimiter:sqlCREATE STREAM csvFileDemo () WITH ( FORMAT = "DELIMITED", DATASOURCE = "abc.csv", TYPE = "file", DELIMITER = " ", CONF_KEY = "csv" );
Tutorial 2: Parse Multi-Line JSON Files
Consider a file with multiple JSON objects separated by newlines:
{"id": 1, "name": "John Doe"}
{"id": 2, "name": "Jane Doe"}
{"id": 3, "name": "John Smith"}Configure
etc/sources/file.yaml:yamljsonlines: fileType: linesDefine a stream with
FORMAT = "JSON":sqlCREATE STREAM linesFileDemo () WITH ( FORMAT = "JSON", TYPE = "file", CONF_KEY = "jsonlines" );
Tutorial 3: Monitor a Directory for New Binary Files
This scenario monitors data/watch for new image files and publishes the raw binary payloads to MQTT.
Step 1: Create the Monitoring Configuration
Create a configuration named watch using the REST API:
PUT http://{{host}}/metadata/sources/file/confKeys/watch
Content-Type: application/json
{
"interval": 0,
"fileType": "raw",
"path": "data"
}interval: 0activates filesystem event notifications.fileType: "raw"reads file content as unparsed binary bytes.path: "data"sets the base directory.
Step 2: Create the Stream
Create the stream using the REST API:
POST http://{{host}}/streams
Content-Type: application/json
{
"sql": "CREATE STREAM watch() WITH (TYPE=\"file\", FORMAT=\"binary\", DATASOURCE=\"watch\", CONF_KEY=\"watch\", SHARED=\"true\");"
}The engine monitors the combined path data/watch.
Step 3: Create the Ingestion Rule
Create a rule to forward binary payloads to MQTT:
POST http://{{host}}/rules
Content-Type: application/json
{
"id": "ruleWatch",
"name": "Watch image folder and send raw binary data to MQTT",
"sql": "SELECT self FROM watch",
"actions": [
{
"mqtt": {
"server": "tcp://127.0.0.1:1883",
"topic": "result",
"sendSingle": true,
"format": "binary"
}
}
]
}Step 4: Validate File Monitoring
Subscribe to MQTT topic result. Copy an image file into data/watch. The engine reads the image file and publishes the binary data to the MQTT broker.