如何在Kafka Connect FilePulse中跳过文件头部元数据行?
Kafka Connect FilePulse 跳过文件头部行配置问题
我正在使用Kafka Connect FilePulse读取文件数据,需要跳过文件头部的多行元数据。示例JSON文件内容如下:
Metadata1 Metadata2 { a: 2 }
我需要跳过前2行,但配置中设置skip.headers = 2后并未生效,当前配置如下:
{ "name": "test", "config": { "connector.class": "io.streamthoughts.kafka.connect.filepulse.source.FilePulseSourceConnector", "filters": "ParseJSON", "filters.ParseJSON.type":"io.streamthoughts.kafka.connect.filepulse.filter.JSONFilter", "filters.ParseJSON.source":"message", "filters.ParseJSON.merge":"true", "filters.ParseJSON.explode.array":"true", "skip.headers": "2", "fs.cleanup.policy.class": "io.streamthoughts.kafka.connect.filepulse.fs.clean.LogCleanupPolicy", "fs.cleanup.policy.triggered.on":"COMMITTED", "fs.listing.class": "io.streamthoughts.kafka.connect.filepulse.fs.LocalFSDirectoryListing", "fs.listing.directory.path":"/tmp/kafka-connect/examples/", "fs.listing.filters":"io.streamthoughts.kafka.connect.filepulse.fs.filter.RegexFileListFilter", "fs.listing.interval.ms": "10000", "file.filter.regex.pattern":".*\\.json$", "topic": "json", "tasks.reader.class": "io.streamthoughts.kafka.connect.filepulse.fs.reader.LocalBytesArrayInputReader", "tasks.file.status.storage.class": "io.streamthoughts.kafka.connect.filepulse.state.KafkaFileObjectStateBackingStore", "tasks.file.status.storage.bootstrap.servers": "broker:29092", "tasks.file.status.storage.topic": "connect-file-pulse-status-1", "tasks.file.status.storage.topic.partitions": 10, "tasks.file.status.storage.topic.replication.factor": 1, "tasks.max": 1 } }
请问该如何正确配置以跳过文件头部的前2行?
解决方案
问题出在你使用的LocalBytesArrayInputReader文件读取器上:skip.headers配置仅对行读取器(如LocalRowFileInputReader)生效,字节数组读取器不支持该参数。
你可以通过以下两种方式解决:
方法一:改用行读取器
将tasks.reader.class替换为行读取器,保留skip.headers配置即可:
"tasks.reader.class": "io.streamthoughts.kafka.connect.filepulse.fs.reader.LocalRowFileInputReader", "skip.headers": "2"
方法二:添加SkipHeaderFilter过滤头部行
如果必须使用字节数组读取器,可通过添加过滤器实现跳过指定行数:
- 在
filters中加入SkipHeader(多过滤器用逗号分隔):
"filters": "SkipHeader,ParseJSON"
- 配置SkipHeaderFilter参数:
"filters.SkipHeader.type": "io.streamthoughts.kafka.connect.filepulse.filter.SkipHeaderFilter", "filters.SkipHeader.lines": "2"
修改后的完整配置示例:
{ "name": "test", "config": { "connector.class": "io.streamthoughts.kafka.connect.filepulse.source.FilePulseSourceConnector", "filters": "SkipHeader,ParseJSON", "filters.SkipHeader.type": "io.streamthoughts.kafka.connect.filepulse.filter.SkipHeaderFilter", "filters.SkipHeader.lines": "2", "filters.ParseJSON.type":"io.streamthoughts.kafka.connect.filepulse.filter.JSONFilter", "filters.ParseJSON.source":"message", "filters.ParseJSON.merge":"true", "filters.ParseJSON.explode.array":"true", "fs.cleanup.policy.class": "io.streamthoughts.kafka.connect.filepulse.fs.clean.LogCleanupPolicy", "fs.cleanup.policy.triggered.on":"COMMITTED", "fs.listing.class": "io.streamthoughts.kafka.connect.filepulse.fs.LocalFSDirectoryListing", "fs.listing.directory.path":"/tmp/kafka-connect/examples/", "fs.listing.filters":"io.streamthoughts.kafka.connect.filepulse.fs.filter.RegexFileListFilter", "fs.listing.interval.ms": "10000", "file.filter.regex.pattern":".*\\.json$", "topic": "json", "tasks.reader.class": "io.streamthoughts.kafka.connect.filepulse.fs.reader.LocalBytesArrayInputReader", "tasks.file.status.storage.class": "io.streamthoughts.kafka.connect.filepulse.state.KafkaFileObjectStateBackingStore", "tasks.file.status.storage.bootstrap.servers": "broker:29092", "tasks.file.status.storage.topic": "connect-file-pulse-status-1", "tasks.file.status.storage.topic.partitions": 10, "tasks.file.status.storage.topic.replication.factor": 1, "tasks.max": 1 } }
内容的提问来源于stack exchange,提问作者prosecutor
相关产品推荐
相关产品推荐

