【问题标题】:Kafka redshift connector throwing NullPointerException卡夫卡红移连接器抛出 NullPointerException
【发布时间】:2020-05-22 00:15:26
【问题描述】:

我正在尝试使用 Redshift Sink Connector 将 Redshift 连接到 Kafka 集群。我目前正在使用独立模式进行测试。

但是,我不断收到以下错误。我相信配置文件是正确的,因为 JdbcDbWriter 连接成功。我还检查了来自 Kafka 的消息不是 Null。我还有更多要检查的吗?非常感谢您。

[2020-05-18 08:08:30,915] INFO Attempting to open connection #1 to Redshift (io.confluent.connect.aws.redshift.jdbc.util.CachedConnectionProvider:87)
[2020-05-18 08:08:31,301] INFO JdbcDbWriter Connected (io.confluent.connect.aws.redshift.jdbc.sink.JdbcDbWriter:49)
[2020-05-18 08:08:31,458] ERROR WorkerSinkTask{id=redshift-source-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted. Error: null (org.apache.kafka.connect.runtime.WorkerSinkTask:566)
java.lang.NullPointerException
        at io.confluent.connect.aws.redshift.jdbc.sink.BufferedRecords.flush(BufferedRecords.java:174)
        at io.confluent.connect.aws.redshift.jdbc.sink.JdbcDbWriter.write(JdbcDbWriter.java:72)
        at io.confluent.connect.aws.redshift.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:74)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:546)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:326)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:228)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:196)
        at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:184)
        at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:234)
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
        at java.util.concurrent.FutureTask.run(FutureTask.java:266)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)
[2020-05-18 08:08:31,467] ERROR WorkerSinkTask{id=redshift-source-0} Task threw an uncaught and unrecoverable exception (org.apache.kafka.connect.runtime.WorkerTask:186)
org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception.
        at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:568)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:326)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:228)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:196)
        at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:184)
        at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:234)
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
        at java.util.concurrent.FutureTask.run(FutureTask.java:266)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)
Caused by: java.lang.NullPointerException
        at io.confluent.connect.aws.redshift.jdbc.sink.BufferedRecords.flush(BufferedRecords.java:174)
        at io.confluent.connect.aws.redshift.jdbc.sink.JdbcDbWriter.write(JdbcDbWriter.java:72)
        at io.confluent.connect.aws.redshift.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:74)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:546)
        ... 10 more
[2020-05-18 08:08:31,469] ERROR WorkerSinkTask{id=redshift-source-0} Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask:187)
[2020-05-18 08:08:31,469] INFO Stopping task (io.confluent.connect.aws.redshift.jdbc.sink.JdbcSinkTask:105)
[2020-05-18 08:08:31,469] INFO Closing connection #1 to Redshift (io.confluent.connect.aws.redshift.jdbc.util.CachedConnectionProvider:113)
[2020-05-18 08:08:31,470] INFO [Consumer clientId=connector-consumer-redshift-source-0, groupId=connect-redshift-source] Revoke previously assigned partitions TopicTest-0 (org.apache.kafka.clients.consumer.internals.ConsumerCoordinator:292)
[2020-05-18 08:08:31,471] INFO [Consumer clientId=connector-consumer-redshift-source-0, groupId=connect-redshift-source] Member connector-consumer-redshift-source-0-c72b053d-8ce9-4337-9757-04e438dd6d0f sending LeaveGroup request to coordinator b-3.inside-dev-messag.ktn4r3.c3.kafka.ap-northeast-1.amazonaws.com:9092 (id: 2147483644 rack: null) due to the consumer is being closed (org.apache.kafka.clients.consumer.internals.AbstractCoordinator:979)

我给了2个属性文件,下面也列出来了。

name=redshift-source
connector.class=io.confluent.connect.aws.redshift.RedshiftSinkConnector
confluent.topic.bootstrap.servers=<Bootstrap Servers>
confluent.topic.replication.factor=2
tasks.max=1

topics=TopicTest

aws.redshift.domain=<Domain>
aws.redshift.port=5439
aws.redshift.database=testdb
aws.redshift.user=<User>
aws.redshift.password=<Pwd>
auto.create=false

type=sink

bootstrap.servers=<Bootstrap Servers>

key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter.schemas.enable=false

internal.key.converter=org.apache.kafka.connect.json.JsonConverter
internal.value.converter=org.apache.kafka.connect.json.JsonConverter
internal.key.converter.schemas.enable=false
internal.value.converter.schemas.enable=false

offset.storage.file.filename=/tmp/connect.offsets
plugin.path=share/java,share/confluent-hub-components

【问题讨论】:

    标签: apache-kafka amazon-redshift apache-kafka-connect


    【解决方案1】:

    我遇到了类似的问题,我的根本原因是未启用架构。在确保消息确实在记录上具有架构后,尝试将以下配置值更改为 true。

    key.converter.schemas.enable=true
    value.converter.schemas.enable=true
    

    要查看消息是否启用了架构,请检查 Kafka 中的 json 消息,您应该会看到“架构”部分和“有效负载”部分。这将告诉您架构是消息中数据的一部分。

    如果您在 json 中没有看到架构,请确保在源连接器上也启用了架构。

    我在 github 上得到的答案帮助了我,你可以看到 here

    另一个相关问题是this

    【讨论】:

    • 嗨!问题是 key 为 null,而问题正是 nullpointexception。我不知道来自 topic 的 kafka 消息的默认键为 null。
    • @PiljaeChae,你是如何解决这个问题的?我和你的情况一样
    猜你喜欢
    • 2020-08-27
    • 2016-09-03
    • 2016-07-26
    • 1970-01-01
    • 1970-01-01
    • 2018-06-23
    • 2017-11-14
    • 1970-01-01
    • 2023-02-15
    相关资源
    最近更新 更多