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

Confluent通用S3源连接器无任务创建问题排查求助

问题

在Confluent平台运行Generalized S3 Source Connector时遇到异常,连接器任务未创建,Connect控制台无额外错误信息,仅SSH控制台输出以下日志:

[2023-02-11 11:12:45,464] INFO [Worker clientId=connect-1, groupId=connect-cluster-1] Finished starting connectors and tasks (org.apache.kafka.connect.runtime.distributed.DistributedHerder:1709)
log4j:ERROR A "io.confluent.log4j.redactor.RedactorAppender" object is not assignable to a "org.apache.log4j.Appender" variable.
log4j:ERROR The class "org.apache.log4j.Appender" was loaded by
log4j:ERROR [PluginClassLoader{pluginLocation=file:/usr/share/java/source-2.5.1/}] whereas object of type
log4j:ERROR "io.confluent.log4j.redactor.RedactorAppender" was loaded by [jdk.internal.loader.ClassLoaders$AppClassLoader@251a69d7].
log4j:ERROR Could not instantiate appender named "redactor".
[2023-02-11 11:13:35,741] INFO Injecting Confluent license properties into connector '<unspecified>' (org.apache.kafka.connect.runtime.WorkerConfigDecorator:412)
[2023-02-11 11:13:44,001] INFO Injecting Confluent license properties into connector 'S3GenConnectorConnector_7' (org.apache.kafka.connect.runtime.WorkerConfigDecorator:412)
[2023-02-11 11:13:44,006] INFO S3SourceConnectorConfig values:
        aws.access.key.id = <<ACCESS KEY HERE>>
        aws.secret.access.key = [hidden]
        behavior.on.error = fail
        bucket.listing.max.objects.threshold = -1
        confluent.license = [hidden]
        confluent.topic = _confluent-command
        confluent.topic.bootstrap.servers = [172.27.157.66:9092]
        confluent.topic.replication.factor = 3
        directory.delim = /
        file.discovery.starting.timestamp = 0
        filename.regex = (.+)\+(\d+)\+.+$\n        folders = []
        format.bytearray.extension = .bin
        format.bytearray.separator =
        format.class = class io.confluent.connect.s3.format.string.StringFormat
        format.json.schema.enable = false
        mode = RESTORE_BACKUP
        parse.error.topic.prefix = error
        partition.field.name = []
        partitioner.class = class io.confluent.connect.storage.partitioner.DefaultPartitioner
        path.format =
        record.batch.max.size = 200
        s3.bucket.name = mytestbucketamtk
        s3.credentials.provider.class = class com.amazonaws.auth.DefaultAWSCredentialsProviderChain
        s3.http.send.expect.continue = true
        s3.part.retries = 3
        s3.path.style.access = true
        s3.poll.interval.ms = 60000
        s3.proxy.password = null
        s3.proxy.url =
        s3.proxy.username = null
        s3.region = us-east-1
        s3.retry.backoff.ms = 200
        s3.sse.customer.key = null
        s3.ssea.name =
        s3.wan.mode = false
        schema.cache.size = 50
        store.url = null
        task.batch.size = 10
        topic.regex.list = [first_topic:.*]
        topics.dir = topics
 (io.confluent.connect.s3.source.S3SourceConnectorConfig:376)
[2023-02-11 11:13:44,029] INFO Using configured AWS access key credentials instead of configured credentials provider class. (io.confluent.connect.s3.source.S3Storage:500)

连接器配置文件内容:

{
  "name": "S3GenConnectorConnector_7",
  "config": {
    "connector.class": "io.confluent.connect.s3.source.S3SourceConnector",
    "tasks.max": "1",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "mode": "RESTORE_BACKUP",
    "format.class": "io.confluent.connect.s3.format.string.StringFormat",
    "s3.bucket.name": "mytestbucketamtk",
    "s3.region": "us-east-1",
    "aws.access.key.id": <<ACCESS KEY HERE>>,
    "aws.secret.access.key": <<SECRET KEY>>,
    "topic.regex.list":"first_topic:.*"
  }
}

当前状况:连接器任务未创建,Connect控制台无其他错误提示,需排查问题根源及是否遗漏必填配置。

排查方向与解决建议

1. 优先处理Log4j类加载冲突问题

日志中明确出现log4j:ERROR Could not instantiate appender named "redactor",原因是RedactorAppender和org.apache.log4j.Appender由不同类加载器加载,这会导致日志输出异常,可能掩盖连接器的真实错误信息:

  • 检查Connect集群的connect-log4j.properties配置文件,确认redactor appender的配置是否正确
  • 确保confluent-log4j-redactor包版本与Connect集群的Log4j版本兼容,避免类加载冲突
  • 临时注释掉redactor appender配置,重启Connect集群后查看是否有更多错误日志输出

2. 检查连接器配置的完整性与匹配性

从日志中的S3SourceConnectorConfig来看,部分配置使用默认值,需重点确认:

  • topics.dir:默认值为topics,需确认S3桶中存在对应路径(s3://mytestbucketamtk/topics/),且路径下有符合topic.regex.list规则的文件
  • filename.regex:默认值为(.+)\+(\d+)\+.+$,需确认S3中的文件名是否匹配该格式(例如first_topic+0000000000+0000000000.json),格式不符会导致连接器无法识别文件
  • AWS权限:虽然日志显示已使用配置的密钥,但需确认该密钥对目标S3桶拥有读取、列出内容的权限

3. 获取更详细的连接器状态信息

  • 查看Connect主日志文件(通常在/var/log/confluent/connect/目录下),查找连接器初始化阶段的任务创建失败细节
  • 使用Confluent CLI命令查看连接器状态:confluent connect connector describe S3GenConnectorConnector_7,获取官方工具输出的状态详情

4. 验证RESTORE_BACKUP模式的适配性

当前使用mode=RESTORE_BACKUP模式,需满足:

  • S3中的文件必须是由Confluent S3 Sink Connector生成的标准格式,否则连接器无法解析
  • 如果是自定义格式文件,需确认format.class配置的格式类能正确解析目标文件内容

内容的提问来源于stack exchange,提问作者Don Woodward

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 12:15:49