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

Kafka-Snowflake连接器报错:snowflake.schema.name schema不存在求助

问题:Snowflake-Kafka连接器启动报错“snowflake.schema.name schema does not exist”

环境信息

  • Kafka版本:kafka_2.13-3.5.0
  • Snowflake-Kafka连接器版本:snowflake-kafka-connector-2.0.0.jar
  • Snowflake账号角色:ACCOUNTADMIN(已配置全量权限)

连接器配置(SF_connect.properties)

name=KafkaConnector
connector.class=com.snowflake.kafka.connector.SnowflakeSinkConnector
tasks.max=8
topics=tokenize-output
buffer.count.records=10000
buffer.flush.time=60
buffer.size.bytes=5000000
snowflake.url.name=<host>.snowflakecomputing.com:443
snowflake.user.name=KAFKA_CONNECTOR_WJ
snowflake.private.key=-----BEGIN RSA PRIVATE KEY----- MII...-----END RSA PRIVATE KEY-----
snowflake.database.name=KAFKA_WJ_TEST
snowflake.schema.name=PUBLIC
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.storage.StringConverter

启动命令

/opt/kafka/kafka_2.13-3.5.0/bin/connect-standalone.sh /opt/kafka/kafka_2.13-3.5.0/config/connect-standalone.properties /opt/kafka/kafka_2.13-3.5.0/config/SF_connect.properties

报错堆栈

[2023-08-17 09:53:03,270] INFO Kafka Connect started (org.apache.kafka.connect.runtime.Connect:57)
[2023-08-17 09:53:03,496] INFO [SF_KAFKA_CONNECTOR] Establishing a JDBC connection with url:jdbc:snowflake://mya39555.snowflakecomputing.com:443 (com.snowflake.kafka.connector.internal.SnowflakeConnectionServiceV1:46)
[2023-08-17 09:53:05,315] INFO [SF_KAFKA_CONNECTOR] initialized the snowflake connection (com.snowflake.kafka.connector.internal.SnowflakeConnectionServiceV1:46)
[2023-08-17 09:53:06,722] INFO [SF_KAFKA_CONNECTOR] database KAFKA_WJ_TEST exists (com.snowflake.kafka.connector.internal.SnowflakeConnectionServiceV1:46)
[2023-08-17 09:53:07,020] ERROR [SF_KAFKA_CONNECTOR] Validate Error msg:[SF_KAFKA_CONNECTOR] Exception: Failed to prepare SQL statement
Error Code: 2001
Detail: SQL Exception, reported by Snowflake JDBC
Message: SQL compilation error:
Object does not exist, or operation cannot be performed.
net.snowflake.client.jdbc.SnowflakeUtil.checkErrorAndThrowExceptionSub(SnowflakeUtil.java:129)
net.snowflake.client.jdbc.SnowflakeUtil.checkErrorAndThrowException(SnowflakeUtil.java:70)
net.snowflake.client.core.StmtUtil.pollForOutput(StmtUtil.java:456)
net.snowflake.client.core.StmtUtil.execute(StmtUtil.java:366)
net.snowflake.client.core.SFStatement.executeHelper(SFStatement.java:464)
net.snowflake.client.core.SFStatement.executeQueryInternal(SFStatement.java:192)
net.snowflake.client.core.SFStatement.executeQuery(SFStatement.java:129)
net.snowflake.client.core.SFStatement.execute(SFStatement.java:742)
net.snowflake.client.core.SFStatement.execute(SFStatement.java:652)
net.snowflake.client.jdbc.SnowflakeStatementV1.executeInternal(SnowflakeStatementV1.java:310)
net.snowflake.client.jdbc.SnowflakePreparedStatementV1.execute(SnowflakePreparedStatementV1.java:458)
com.snowflake.kafka.connector.internal.SnowflakeConnectionServiceV1.schemaExists(SnowflakeConnectionServiceV1.java:655)
com.snowflake.kafka.connector.SnowflakeSinkConnector.validate(SnowflakeSinkConnector.java:270)
org.apache.kafka.connect.runtime.AbstractHerder.validateConnectorConfig(AbstractHerder.java:509)
org.apache.kafka.connect.runtime.AbstractHerder.lambda$validateConnectorConfig$2(AbstractHerder.java:390)
java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
java.base/java.lang.Thread.run(Thread.java:829), errorCode:2001 (com.snowflake.kafka.connector.SnowflakeSinkConnector:94)
[2023-08-17 09:53:07,022] INFO AbstractConfig values:
 (org.apache.kafka.common.config.AbstractConfig:369)
[2023-08-17 09:53:07,037] ERROR Failed to create connector for /opt/kafka/kafka_2.13-3.5.0/config/SF_connect.properties (org.apache.kafka.connect.cli.ConnectStandalone:74)
[2023-08-17 09:53:07,038] ERROR Stopping after connector error (org.apache.kafka.connect.cli.ConnectStandalone:84)
java.util.concurrent.ExecutionException: org.apache.kafka.connect.runtime.rest.errors.BadRequestException: Connector configuration is invalid and contains the following 1 error(s):
snowflake.schema.name schema does not exist
You can also find the above list of errors at the endpoint `/connector-plugins/{connectorType}/config/validate`
        at org.apache.kafka.connect.util.ConvertingFutureCallback.result(ConvertingFutureCallback.java:123)
        at org.apache.kafka.connect.util.ConvertingFutureCallback.get(ConvertingFutureCallback.java:107)
        at org.apache.kafka.connect.cli.ConnectStandalone.processExtraArgs(ConnectStandalone.java:81)
        at org.apache.kafka.connect.cli.AbstractConnectCli.startConnect(AbstractConnectCli.java:150)
        at org.apache.kafka.connect.cli.AbstractConnectCli.run(AbstractConnectCli.java:94)
        at org.apache.kafka.connect.cli.ConnectStandalone.main(ConnectStandalone.java:112)
Caused by: org.apache.kafka.connect.runtime.rest.errors.BadRequestException: Connector configuration is invalid and contains the following 1 error(s):
snowflake.schema.name schema does not exist
You can also find the above list of errors at the endpoint `/connector-plugins/{connectorType}/config/validate`
        at org.apache.kafka.connect.runtime.AbstractHerder.maybeAddConfigErrors(AbstractHerder.java:754)
        at org.apache.kafka.connect.runtime.standalone.StandaloneHerder.putConnectorConfig(StandaloneHerder.java:207)
        at org.apache.kafka.connect.runtime.standalone.StandaloneHerder.lambda$null$0(StandaloneHerder.java:193)
        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.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:304)
        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)
[2023-08-17 09:53:07,039] INFO Kafka Connect stopping (org.apache.kafka.connect.runtime.Connect:67)

已做排查动作

  • 尝试过使用Snowflake默认的PUBLIC schema和自定义创建的schema,均报错
  • 已在Snowflake中创建对应Topic的表(有无表均触发相同报错)
  • 确认账号拥有ACCOUNTADMIN角色及全量权限
  • 日志显示数据库KAFKA_WJ_TEST已成功识别存在

解决建议

1. 修正私钥格式

配置中的snowflake.private.key存在空格(-----BEGIN RSA PRIVATE KEY----- MII...),这会导致私钥解析异常。正确做法是:

  • 如果私钥是多行格式,在properties文件中用反斜杠\转义每一行的换行
  • 或者将私钥合并为一行,去掉所有空格和换行,确保从-----BEGIN RSA PRIVATE KEY-----到-----END RSA PRIVATE KEY-----是连续的字符串

2. 显式指定Snowflake角色

在配置文件中添加:

snowflake.role.name=ACCOUNTADMIN

即使账号默认角色是ACCOUNTADMIN,连接器会话可能未正确继承,显式指定可确保会话使用该高权限角色。

3. 添加仓库配置

Snowflake会话需要绑定仓库才能访问数据库对象,在配置中添加:

snowflake.warehouse.name=<你的仓库名称>

替换为你实际使用的Snowflake仓库名称。

4. 验证大小写匹配

检查Snowflake中KAFKA_WJ_TEST数据库和PUBLICschema的创建方式:

  • 如果创建时使用了双引号(如CREATE DATABASE "KAFKA_WJ_TEST"),则配置中的名称必须完全匹配大小写
  • 默认创建的对象是大小写不敏感的,配置中使用大写或小写均可,但建议保持一致

5. 手动测试JDBC连接

用相同的账号、私钥,通过Snowflake JDBC客户端执行以下命令,验证是否能正常访问schema:

USE DATABASE KAFKA_WJ_TEST;
USE SCHEMA PUBLIC;
SHOW SCHEMAS IN DATABASE KAFKA_WJ_TEST;

如果命令执行失败,说明存在权限或对象存在性问题,需优先解决。

6. 升级连接器版本

snowflake-kafka-connector-2.0.0与kafka_2.13-3.5.0可能存在兼容性问题,建议升级到最新稳定版本(如2.8.x及以上),官方文档中已明确高版本连接器对新Kafka版本的支持。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 04:23:16