如何在Kafka Connect的kafka-connect-spooldir中解析dd.MM.yyyy格式日期?
核心问题
你的配置中parser.timestamp.date.formats的JSON语法有误,导致连接器没有正确加载日期解析规则,进而无法匹配dd.MM.yyyy格式的日期字符串。另外,需确保Value Schema的日期类型定义与解析规则正确对应。
修复步骤
修正日期解析格式的配置
在JSON格式的连接器配置中,parser.timestamp.date.formats的每个格式字符串必须用双引号包裹,数组结构要严格符合JSON语法。你的原配置中缺少字符串引号,导致连接器无法识别预设的日期格式。正确配置应为:
"parser.timestamp.date.formats": ["dd.MM.yyyy", "yyyy-MM-dd'T'HH:mm:ss", "yyyy-MM-dd' 'HH:mm:ss"]确认Value Schema的日期类型定义
对于dd.MM.yyyy格式的纯日期(无时间部分),使用org.apache.kafka.connect.data.Date类型是合适的,它对应INT32类型(存储自1970-01-01以来的天数)。你的Schema定义本身是正确的,无需修改,但要确保与修正后的日期解析格式配合使用。如果后续需要解析带时间的日期字符串,切换为
org.apache.kafka.connect.data.Timestamp类型(对应INT64,存储毫秒数)即可,同样会匹配你配置的日期格式。额外补充:小数分隔符处理
你的CSV中小数用逗号(如-55,60)作为分隔符,需要添加csv.decimal.separator.char=44(逗号的ASCII码)配置,避免Umsatz字段解析失败。
完整的连接器配置示例
{ "connector.class": "com.github.jcustenborder.kafka.connect.spooldir.SpoolDirCsvSourceConnector", "csv.first.row.as.header": "true", "finished.path": "/csv/finished", "tasks.max": "1", "parser.timestamp.date.formats": ["dd.MM.yyyy", "yyyy-MM-dd'T'HH:mm:ss", "yyyy-MM-dd' 'HH:mm:ss"], "key.schema": {"name":"com.github.jcustenborder.kafka.connect.model.Key","type":"STRUCT","isOptional":false,"fieldSchemas":{}}, "csv.separator.char": "59", "input.file.pattern": "umsaetze_.*.csv", "topic": "test-csv", "error.path": "/csv/error", "input.path": "/csv/input", "csv.decimal.separator.char": "44", "value.schema": { "name": "com.github.jcustenborder.kafka.connect.model.Value", "type": "STRUCT", "isOptional": false, "fieldSchemas": { "Buchungstag": { "name": "org.apache.kafka.connect.data.Date", "type": "INT32", "version": 1, "isOptional": true }, "Wertstellung": { "name": "org.apache.kafka.connect.data.Date", "type": "INT32", "version": 1, "isOptional": true }, "Vorgang": { "type": "STRING", "isOptional": false }, "Buchungstext": { "type": "STRING", "isOptional": false }, "Umsatz": { "name": "org.apache.kafka.connect.data.Decimal", "type": "BYTES", "version": 1, "parameters": { "scale": "2" }, "isOptional": true } } } }
验证逻辑
修正后,连接器会优先使用dd.MM.yyyy格式解析Buchungstag和Wertstellung字段的日期字符串,匹配成功后转换为对应的Date类型存储到Kafka中。如果遇到其他格式的日期字符串(如带时间戳的),会自动尝试数组中的后续格式进行解析。
内容的提问来源于stack exchange,提问作者Regfor

