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
相关产品推荐
相关产品推荐

