You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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的分区子目录,若不能则提前创建或赋予递归创建目录的权限。

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 00:00:03