【问题标题】:Can varchar datatype be a timestamp in Confluent?varchar 数据类型可以是 Confluent 中的时间戳吗?
【发布时间】:2018-12-25 14:56:28
【问题描述】:

我正在使用 Confluent 来实现实时 ETL。 我的数据源是oracle,每个表都有一个名为ts的列,它的数据类型是varchar,但是该列的数据是YYYY-MM--DD HH24:MI:SS格式。 我可以将此列用作融合 kafka 连接器中的时间戳吗? 如何配置 xxxxx.properties 文件?

mode=timestamp
query= select to_date(a.ts,'yyyy-mm-dd hh24:mi:ss') tsinc,a.* from TEST_CORP a
poll.interval.ms=1000 
timestamp.column.name=tsinc

【问题讨论】:

  • 尝试查看 TimestampConverter 的单个消息转换

标签: oracle jdbc apache-kafka apache-kafka-connect confluent-platform


【解决方案1】:

connector.class=io.confluent.connect.jdbc.JdbcSourceConnector 查询=从 NFSN.BD_CORP 中选择 * 模式=时间戳 poll.interval.ms=3000 timestamp.column.name=TS topic.prefix=t_validate.non.null=false

然后我得到这个错误:

[2018-12-25 14:39:59,756] INFO 过滤后的表格是: (io.confluent.connect.jdbc.source.TableMonitorThread:175) [2018-12-25 14:40:01,383] 调试检查下一个结果块 TimestampIncrementingTableQuerier{table=null, query='select * from NFSN.BD_CORP',topicPrefix='t_',incrementingColumn='', 时间戳列=[TS]} (io.confluent.connect.jdbc.source.JdbcSourceTask:291) [2018-12-25 14:40:01,386] 调试 TimestampIncrementingTableQuerier{table=null, query='select * from NFSN.BD_CORP', topicPrefix='t_', incrementingColumn='', timestampColumns=[TS]} 准备好的 SQL 查询: 选择 * 从 NFSN.BD_CORP WHERE "TS" > ?和 "TS"

    at oracle.jdbc.driver.T4CTTIoer.processError(T4CTTIoer.java:447)
    at oracle.jdbc.driver.T4CTTIoer.processError(T4CTTIoer.java:396)
    at oracle.jdbc.driver.T4C8Oall.processError(T4C8Oall.java:951)
    at oracle.jdbc.driver.T4CTTIfun.receive(T4CTTIfun.java:513)
    at oracle.jdbc.driver.T4CTTIfun.doRPC(T4CTTIfun.java:227)
    at oracle.jdbc.driver.T4C8Oall.doOALL(T4C8Oall.java:531)
    at oracle.jdbc.driver.T4CPreparedStatement.doOall8(T4CPreparedStatement.java:208)
    at oracle.jdbc.driver.T4CPreparedStatement.executeForDescribe(T4CPreparedStatement.java:886)
    at oracle.jdbc.driver.OracleStatement.executeMaybeDescribe(OracleStatement.java:1175)
    at oracle.jdbc.driver.OracleStatement.doExecuteWithTimeout(OracleStatement.java:1296)
    at oracle.jdbc.driver.OraclePreparedStatement.executeInternal(OraclePreparedStatement.java:3613)
    at oracle.jdbc.driver.OraclePreparedStatement.executeQuery(OraclePreparedStatement.java:3657)
    at oracle.jdbc.driver.OraclePreparedStatementWrapper.executeQuery(OraclePreparedStatementWrapper.java:1495)
    at io.confluent.connect.jdbc.source.TimestampIncrementingTableQuerier.executeQuery(TimestampIncrementingTableQuerier.java:168)
    at io.confluent.connect.jdbc.source.TableQuerier.maybeStartQuery(TableQuerier.java:88)
    at io.confluent.connect.jdbc.source.TimestampIncrementingTableQuerier.maybeStartQuery(TimestampIncrementingTableQuerier.java:60)
    at io.confluent.connect.jdbc.source.JdbcSourceTask.poll(JdbcSourceTask.java:292)
    at org.apache.kafka.connect.runtime.WorkerSourceTask.poll(WorkerSourceTask.java:244)
    at org.apache.kafka.connect.runtime.WorkerSourceTask.execute(WorkerSourceTask.java:220)
    at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:175)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:219)
    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) [2018-12-25 14:40:01,390] DEBUG Resetting querier

TimestampIncrementingTableQuerier{table=null, query='select * from NFSN.BD_CORP',topicPrefix='t_',incrementingColumn='', 时间戳列=[TS]} (io.confluent.connect.jdbc.source.JdbcSourceTask:332) ^C[2018-12-25 14:40:03,826] 信息卡夫卡连接停止 (org.apache.kafka.connect.runtime.Connect:65) [2018-12-25 14:40:03,827] 信息停止 REST 服务器 (org.apache.kafka.connect.runtime.rest.RestServer:223)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2014-02-22
    • 2021-08-15
    • 1970-01-01
    • 2010-11-16
    • 1970-01-01
    • 2018-06-24
    • 2012-02-17
    • 1970-01-01
    相关资源
    最近更新 更多