Confluent SFTP Sink Connector报错:无法通过现有会话打开新SFTP通道
问题描述
使用Confluent SFTP Sink Connector时,每次尝试向SFTP服务器写入数据都会抛出Failed to open new SFTP channel with existing session错误,但用相同凭据可以正常连接SFTP服务器。
错误堆栈信息
org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception. at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:618) at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:334) at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:235) at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:204) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:201) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:256) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: org.apache.kafka.connect.errors.ConnectException: Failed to open new SFTP channel with existing session. at io.confluent.connect.sftp.connection.SftpConnection.newChannelFromSession(SftpConnection.java:108) at io.confluent.connect.sftp.sink.storage.SftpOutputStream.createPath(SftpOutputStream.java:48) at io.confluent.connect.sftp.sink.storage.SftpOutputStream.<init>(SftpOutputStream.java:44) at io.confluent.connect.sftp.sink.storage.SftpSinkStorage.create(SftpSinkStorage.java:87) at io.confluent.connect.sftp.sink.format.csv.CsvRecordWriter.<init>(CsvRecordWriter.java:32) at io.confluent.connect.sftp.sink.format.csv.CsvRecordWriterProvider.getRecordWriter(CsvRecordWriterProvider.java:29) at io.confluent.connect.sftp.sink.format.csv.CsvRecordWriterProvider.getRecordWriter(CsvRecordWriterProvider.java:12) at io.confluent.connect.sftp.sink.TopicPartitionWriter.getWriter(TopicPartitionWriter.java:408) at io.confluent.connect.sftp.sink.TopicPartitionWriter.writeRecord(TopicPartitionWriter.java:454) at io.confluent.connect.sftp.sink.TopicPartitionWriter.checkRotationOrAppend(TopicPartitionWriter.java:250) at io.confluent.connect.sftp.sink.TopicPartitionWriter.writePartitionWhe..
当前连接器配置
apiVersion: platform.confluent.io/v1beta1 kind: Connector metadata: name: update-sftp-sink-connector namespace: confluent spec: name: update-sftp-sink-connector taskMax: 1 class: io.confluent.connect.sftp.SftpSinkConnector configs: topics: updates-topic file.delim: "." sftp.host: "my.host.name" sftp.port: "22" sftp.username: ${file:/mnt/secrets/connect-connector-secrets/sftp-creds:username} sftp.password: ${file:/mnt/secrets/connect-connector-secrets/sftp-creds:password} sftp.working.dir: "/updates" directory.delim: "/" rotate.interval.ms: "120000" flush.size: 1 partition.duration.ms: "120000" partitioner.class: io.confluent.connect.storage.partitioner.TimeBasedPartitioner format.class: io.confluent.connect.sftp.sink.format.csv.CsvFormat key.converter: org.apache.kafka.connect.storage.StringConverter value.converter: org.apache.kafka.connect.json.JsonConverter storage.class: io.confluent.connect.sftp.sink.storage.SftpSinkStorage locale: en-GB timezone: UTC timestamp.extractor: Record restartPolicy: type: OnFailure maxRetry: 10 connectClusterRef: name: connect namespace: confluent
排查方向与解决方法
1. SFTP会话超时或被提前关闭
SFTP服务器可能设置了会话空闲超时,连接器复用的会话已被服务器断开,但仍尝试用该会话创建新通道。
- 解决方法:
- 在连接器配置中添加
sftp.session.max.idle.ms参数,设置小于服务器空闲超时的值(例如服务器超时300秒则设为240000),让连接器主动刷新会话。 - 检查SFTP服务器的
ClientAliveInterval和ClientAliveCountMax配置,调整参数避免会话过早断开。
- 在连接器配置中添加
2. 目录/文件操作权限不足
虽能连接SFTP服务器,但连接器可能无目标目录的写入权限,或无法自动创建分区子目录。
- 解决方法:
- 用配置的凭据手动登录SFTP服务器,在
/updates目录下尝试创建文件、写入内容,确认权限是否正常。 - 若使用时间分区器,检查是否能自动生成如
/updates/2024-05-20的分区子目录,若不能则提前创建或赋予递归创建目录的权限。
- 用配置的凭据手动登录SFTP服务器,在
3. 连接器版本兼容性问题
旧版本Confluent SFTP Sink Connector可能存在会话复用的bug。
- 解决方法:
- 将连接器升级至最新稳定版本(建议至少5.5.0以上,该版本修复了多个SFTP会话相关问题)。
- 若升级有阻碍,添加
sftp.connection.max.age.ms参数,强制连接器定期重建会话,避免复用失效会话。
4. 网络层面会话中断
防火墙或代理可能在TCP层面中断空闲SFTP会话,导致连接器持有的会话失效。
- 解决方法:
- 检查网络设备的超时设置,确保SFTP会话的TCP连接不会被中途断开。
- 在连接器配置中添加
sftp.enable.keep.alive=true,启用SFTP保活机制维持会话活跃。
5. 凭据或配置解析问题
配置中的变量引用可能存在解析错误,导致会话创建后权限不足或异常。
- 解决方法:
- 测试环境下临时替换变量为明文凭据,确认是否为变量解析问题。
- 检查secret文件格式,确保
username和password字段存在且无多余空格。
内容的提问来源于stack exchange,提问作者Cthomps
相关产品推荐
相关产品推荐

