Skip to content

AxonOps — AI-Native Control Plane for Open Source Data Platforms

File Source Connector

Stream file contents to Kafka topics for log aggregation and file-based data ingestion.


The File Source Connector reads files and writes their contents to Kafka, supporting:

  • Line-by-line file reading
  • File watching for new content
  • Multiple file patterns
  • Offset tracking for restarts

Production Use

The bundled FileStreamSourceConnector is intended for development and testing. For production file streaming, consider specialized connectors like Filebeat with Kafka output or dedicated log shipping solutions.


Apache Kafka includes a basic file source connector:

{
"name": "file-source",
"config": {
"connector.class": "org.apache.kafka.connect.file.FileStreamSourceConnector",
"tasks.max": "1",
"file": "/var/log/application.log",
"topic": "application-logs"
}
}
LimitationDescription
Single fileOnly reads one file per connector
No glob patternsCannot watch directory
No offset persistenceRestarts from beginning
No rotation handlingDoes not follow rotated logs

For production use, the SpoolDir connector provides advanced file handling:

Terminal window
confluent-hub install jcustenborder/kafka-connect-spooldir:latest
{
"name": "spooldir-source",
"config": {
"connector.class": "com.github.jcustenborder.kafka.connect.spooldir.SpoolDirCsvSourceConnector",
"tasks.max": "1",
"input.path": "/data/incoming",
"finished.path": "/data/processed",
"error.path": "/data/error",
"topic": "file-data",
"input.file.pattern": ".*\\.csv$"
}
}

{
"connector.class": "com.github.jcustenborder.kafka.connect.spooldir.SpoolDirCsvSourceConnector",
"csv.first.row.as.header": "true",
"csv.separator.char": ",",
"csv.quote.char": "\"",
"schema.generation.enabled": "true"
}
{
"connector.class": "com.github.jcustenborder.kafka.connect.spooldir.SpoolDirJsonSourceConnector",
"schema.generation.enabled": "true"
}
{
"connector.class": "com.github.jcustenborder.kafka.connect.spooldir.SpoolDirLineDelimitedSourceConnector"
}
{
"connector.class": "com.github.jcustenborder.kafka.connect.spooldir.SpoolDirBinaryFileSourceConnector"
}

{
"input.path": "/data/incoming",
"input.file.pattern": ".*\\.csv$",
"file.minimum.age.ms": "5000"
}
PropertyDescriptionDefault
input.pathDirectory to watchRequired
input.file.patternRegex for file matching.*
file.minimum.age.msMin file age before processing0
{
"finished.path": "/data/processed",
"error.path": "/data/error",
"cleanup.policy": "MOVE"
}
PolicyDescription
MOVEMove to finished/error path
DELETEDelete after processing
NONELeave in place

{
"schema.generation.enabled": "true",
"schema.generation.key.fields": "id"
}
{
"schema.generation.enabled": "false",
"key.schema": "{\"name\":\"key\",\"type\":\"STRING\"}",
"value.schema": "{\"name\":\"event\",\"type\":\"STRUCT\",\"fields\":[{\"name\":\"id\",\"type\":\"STRING\"},{\"name\":\"timestamp\",\"type\":\"INT64\"},{\"name\":\"data\",\"type\":\"STRING\"}]}"
}

{
"batch.size": "1000",
"empty.poll.wait.ms": "500"
}
PropertyDescriptionDefault
batch.sizeRecords per batch1000
empty.poll.wait.msWait when no files500

{
"error.path": "/data/error",
"halt.on.error": "false"
}
{
"errors.tolerance": "all",
"errors.deadletterqueue.topic.name": "dlq-file-source",
"errors.deadletterqueue.topic.replication.factor": 3,
"errors.log.enable": "true"
}

{
"name": "csv-file-source",
"config": {
"connector.class": "com.github.jcustenborder.kafka.connect.spooldir.SpoolDirCsvSourceConnector",
"tasks.max": "1",
"input.path": "/data/incoming/csv",
"finished.path": "/data/processed/csv",
"error.path": "/data/error/csv",
"input.file.pattern": ".*\\.csv$",
"file.minimum.age.ms": "10000",
"topic": "csv-events",
"csv.first.row.as.header": "true",
"csv.separator.char": ",",
"csv.quote.char": "\"",
"csv.escape.char": "\\",
"csv.null.field.indicator": "NULL",
"schema.generation.enabled": "true",
"schema.generation.key.fields": "id",
"batch.size": "1000",
"cleanup.policy": "MOVE",
"key.converter": "org.apache.kafka.connect.storage.StringConverter",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
"errors.tolerance": "all",
"errors.deadletterqueue.topic.name": "dlq-csv-source",
"errors.deadletterqueue.topic.replication.factor": 3
}
}

{
"name": "jsonl-file-source",
"config": {
"connector.class": "com.github.jcustenborder.kafka.connect.spooldir.SpoolDirJsonSourceConnector",
"tasks.max": "1",
"input.path": "/data/incoming/json",
"finished.path": "/data/processed/json",
"error.path": "/data/error/json",
"input.file.pattern": ".*\\.jsonl$",
"topic": "json-events",
"schema.generation.enabled": "true",
"batch.size": "500",
"cleanup.policy": "MOVE",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false"
}
}

For streaming log files (append-only), consider the FilePulse connector:

Terminal window
confluent-hub install streamthoughts/kafka-connect-file-pulse:latest
{
"name": "log-file-source",
"config": {
"connector.class": "io.streamthoughts.kafka.connect.filepulse.source.FilePulseSourceConnector",
"tasks.max": "1",
"fs.listing.class": "io.streamthoughts.kafka.connect.filepulse.fs.LocalFSDirectoryListing",
"fs.listing.directory.path": "/var/log/app",
"fs.listing.filters": "io.streamthoughts.kafka.connect.filepulse.fs.filter.RegexFileListFilter",
"file.filter.regex.pattern": ".*\\.log$",
"fs.cleanup.policy.class": "io.streamthoughts.kafka.connect.filepulse.fs.clean.LogCleanupPolicy",
"topic": "application-logs",
"tasks.reader.class": "io.streamthoughts.kafka.connect.filepulse.fs.reader.LocalRowFileInputReader",
"offset.strategy": "name+hash"
}
}

{
"tasks.max": "3",
"tasks.file.status.storage.class": "io.streamthoughts.kafka.connect.filepulse.state.InMemoryFileObjectStateBackingStore"
}
{
"batch.size": "5000",
"poll.interval.ms": "1000"
}

Terminal window
curl http://connect:8083/connectors/file-source/status
Terminal window
# Check input directory
ls -la /data/incoming/
# Check processed files
ls -la /data/processed/
# Check error files
ls -la /data/error/

IssueCauseSolution
Files not processedPattern mismatchVerify input.file.pattern
Permission deniedFile permissionsCheck connector user access
Schema errorsInconsistent dataEnable schema generation
Files in error pathParsing failuresCheck file format
High memory usageLarge batch sizeReduce batch.size