Embulk多主键增量同步忽略最后记录及参数绑定错误排查
在使用Embulk从Oracle数据库执行多主键增量数据同步时,遇到以下问题:
- 同步过程忽略最后一条记录
- 任务运行报错:
INまたはOUTパラメータがありません - 索引:: 3(无IN/OUT参数-索引::3) - 出现转换错误:
Converting last_record value 1 to column index 2 is not supported
错误配置案例
exec:<omit> in: type: oracle {% include 'inc/oracle_axlink' %} # connection info query: | SELECT CTOIAWASE_NO, CTOIAWASE_YY, NSYDN_ORDER FROM TEST.TABLE1 WHERE CTOIAWASE_NO >= :CTOIAWASE_NO AND CTOIAWASE_YY >= :CTOIAWASE_YY AND NSYDN_ORDER >= :NSYDN_ORDER use_raw_query_with_incremental: true incremental_columns: [ CTOIAWASE_NO, CTOIAWASE_YY, NSYDN_ORDER ] incremental: true last_record: [ "00100002", "2020", 1 ] options: {characterEncoding: MS932, characterSetResults: MS932,serverTimezone: JST} fetch_rows: 10000 default_column_options: NUMERIC: { value_type: long } DECIMAL: { value_type: long } out:<omit>
错误日志
$ embulk run test.yml.liquid -c test_diff.yml.liquid 2022-10-13 15:51:59.959 +0900: Embulk v0.9.23 2022-10-13 15:52:01.231 +0900 [WARN] (main): DEPRECATION: JRuby org.jruby.embed.ScriptingContainer is directly injected. 2022-10-13 15:52:04.168 +0900 [INFO] (main): Gem's home and path are set by default: "/home/dbdevelop/.embulk/lib/gems" 2022-10-13 15:52:07.121 +0900 [INFO] (main): Started Embulk v0.9.23 2022-10-13 15:52:07.323 +0900 [INFO] (0001:transaction): Loaded plugin embulk-input-oracle (0.9.3) 2022-10-13 15:52:07.428 +0900 [INFO] (0001:transaction): Loaded plugin embulk-output-s3_parquet (0.5.2) 2022-10-13 15:52:07.527 +0900 [INFO] (0001:transaction): Connecting to jdbc:oracle:thin:@XXX.XXXX.XXX.XXX:1521/TEST options {oracle.jdbc.ReadTimeout=1800000, user=XXXXXX, serverTimezone=JST, password=***, characterEncoding=MS932, oracle.net.CONNECT_TIMEOUT=300000, characterSetResults=MS932} 2022-10-13 15:52:08.025 +0900 [INFO] (0001:transaction): Using JDBC Driver 12.1.0.2.0 2022-10-13 15:52:08.109 +0900 [INFO] (0001:transaction): Using local thread executor with max_threads=1 / tasks=1 2022-10-13 15:52:08.565 +0900 [INFO] (0001:transaction): === Output Parquet Schema === 2022-10-13 15:52:08.566 +0900 [INFO] (0001:transaction): message embulk { 2022-10-13 15:52:08.566 +0900 [INFO] (0001:transaction): optional binary CTOIAWASE_NO (STRING); 2022-10-13 15:52:08.566 +0900 [INFO] (0001:transaction): optional binary CTOIAWASE_YY (STRING); 2022-10-13 15:52:08.566 +0900 [INFO] (0001:transaction): optional int64 NSYDN_ORDER; 2022-10-13 15:52:08.566 +0900 [INFO] (0001:transaction): } 2022-10-13 15:52:08.566 +0900 [INFO] (0001:transaction): ============================= 2022-10-13 15:52:08.584 +0900 [INFO] (0001:transaction): {done: 0 / 1, running: 0} 2022-10-13 15:52:08.884 +0900 [WARN] (0017:task-0000): Unable to load native-hadoop library for your platform... using builtin-java classes where applicable log4j:WARN No appenders could be found for logger (org.apache.htrace.core.Tracer). log4j:WARN Please initialize the log4j system properly. log4j:WARN See http://logging.apache.org/log4j/1.2/faq.html#noconfig for more info. 2022-10-13 15:52:09.016 +0900 [INFO] (0017:task-0000): Got brand-new compressor [.snappy] 2022-10-13 15:52:09.347 +0900 [INFO] (0017:task-0000): Local Buffer File: /tmp/embulk-output-s3_parquet-12345678912345/embulk-output-s3_parquet-task-0-0.parquet, Destination: s3://testtesttest-1234567890/append/test/table1.snappy.parquet 2022-10-13 15:52:09.403 +0900 [INFO] (0017:task-0000): Connecting to jdbc:oracle:thin:@XXX.XXX.XXX.XXX:1521/TEST options {oracle.jdbc.ReadTimeout=1800000, user=XXXXXXX, serverTimezone=JST, password=***, characterEncoding=MS932, oracle.net.CONNECT_TIMEOUT=300000, characterSetResults=MS932} 2022-10-13 15:52:09.455 +0900 [INFO] (0017:task-0000): SQL: SELECT CTOIAWASE_NO, CTOIAWASE_YY, NSYDN_ORDER FROM TEST.TABLE1 WHERE CTOIAWASE_NO >= ? AND CTOIAWASE_YY >= :CTOIAWASE_YY AND NSYDN_ORDER >= ? 2022-10-13 15:52:09.455 +0900 [INFO] (0017:task-0000): Parameters: ["2020", 1] 2022-10-13 15:52:09.703 +0900 [INFO] (0001:transaction): {done: 1 / 1, running: 0} org.embulk.exec.PartialExecutionException: java.lang.RuntimeException: java.sql.SQLException: INまたはOUTパラメータがありま せん - 索引:: 3 at org.embulk.exec.BulkLoader$LoaderState.buildPartialExecuteException(BulkLoader.java:340) at org.embulk.exec.BulkLoader.doRun(BulkLoader.java:566) at org.embulk.exec.BulkLoader.access$000(BulkLoader.java:35) at org.embulk.exec.BulkLoader$1.run(BulkLoader.java:353) at org.embulk.exec.BulkLoader$1.run(BulkLoader.java:350) at org.embulk.spi.Exec.doWith(Exec.java:22) at org.embulk.exec.BulkLoader.run(BulkLoader.java:350) at org.embulk.EmbulkEmbed.run(EmbulkEmbed.java:242) at org.embulk.EmbulkRunner.runInternal(EmbulkRunner.java:291) at org.embulk.EmbulkRunner.run(EmbulkRunner.java:155) at org.embulk.cli.EmbulkRun.runSubcommand(EmbulkRun.java:431) at org.embulk.cli.EmbulkRun.run(EmbulkRun.java:90) at org.embulk.cli.Main.main(Main.java:64) Caused by: java.lang.RuntimeException: java.sql.SQLException: INまたはOUTパラメータがありません - 索引:: 3 at com.google.common.base.Throwables.propagate(Throwables.java:160) at org.embulk.input.jdbc.AbstractJdbcInputPlugin.run(AbstractJdbcInputPlugin.java:509) at org.embulk.spi.util.Executors.process(Executors.java:62) at org.embulk.spi.util.Executors.process(Executors.java:38) at org.embulk.exec.LocalExecutorPlugin$DirectExecutor$1.call(LocalExecutorPlugin.java:170) at org.embulk.exec.LocalExecutorPlugin$DirectExecutor$1.call(LocalExecutorPlugin.java:167) 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.sql.SQLException: INまたはOUTパラメータがありません - 索引:: 3 at oracle.jdbc.driver.OraclePreparedStatement.processCompletedBindRow(OraclePreparedStatement.java:2076) at oracle.jdbc.driver.OraclePreparedStatement.executeInternal(OraclePreparedStatement.java:4790) at oracle.jdbc.driver.OraclePreparedStatement.executeQuery(OraclePreparedStatement.java:4845) at oracle.jdbc.driver.OraclePreparedStatementWrapper.executeQuery(OraclePreparedStatementWrapper.java:1501) at org.embulk.input.jdbc.JdbcInputConnection$SingleSelect.fetch(JdbcInputConnection.java:194) at org.embulk.input.jdbc.AbstractJdbcInputPlugin.fetch(AbstractJdbcInputPlugin.java:571) at org.embulk.input.jdbc.AbstractJdbcInputPlugin.run(AbstractJdbcInputPlugin.java:480) ... 8 more Error: java.lang.RuntimeException: java.sql.SQLException: INまたはOUTパラメータがありません - 索引:: 3
Error: java.lang.RuntimeException: java.sql.SQLException: INまたはOUTパラメータがありません - 索引:: 3 is Error: org.embulk.spi.DataException: Converting last_record value 1 to column index 2 is not supported
问题原因分析
参数绑定错误
日志显示生成的SQL混合使用了匿名占位符?和命名占位符:CTOIAWASE_YY,但Embulk在use_raw_query_with_incremental模式下仅支持匿名占位符?,且参数需严格按照incremental_columns的顺序传递。命名占位符未被正确替换,导致参数数量与SQL中的占位符数量不匹配,触发「无IN/OUT参数-索引::3」错误。多主键增量逻辑错误
配置中使用CTOIAWASE_NO >= ? AND CTOIAWASE_YY >= ? AND NSYDN_ORDER >= ?的条件,会错误过滤数据(比如CTOIAWASE_NO更大但CTOIAWASE_YY更小的记录),同时可能导致遗漏最后一条记录。多主键增量需使用组合层级条件,逐次判断主键的大小关系。类型转换与索引不匹配
因参数绑定错误,实际传递的参数缺失第一个值,导致Embulk尝试将1映射到错误的列索引,触发「Converting last_record value 1 to column index 2 is not supported」转换错误。
解决办法
1. 修正SQL参数占位符与增量逻辑
将SQL中的命名占位符改为匿名占位符?,并使用正确的多主键组合条件:
SELECT CTOIAWASE_NO, CTOIAWASE_YY, NSYDN_ORDER FROM TEST.TABLE1 WHERE (CTOIAWASE_NO > ?) OR (CTOIAWASE_NO = ? AND CTOIAWASE_YY > ?) OR (CTOIAWASE_NO = ? AND CTOIAWASE_YY = ? AND NSYDN_ORDER >= ?)
注:多主键的每个层级需要重复前面的主键占位符,确保逻辑正确(大于当前主键,或等于当前主键且下一级主键更大/等于)。
2. 调整incremental_columns与last_record配置
确保incremental_columns的顺序与SQL中占位符的顺序完全对应,last_record的参数数量和类型与占位符匹配:
incremental_columns: [CTOIAWASE_NO, CTOIAWASE_NO, CTOIAWASE_YY, CTOIAWASE_NO, CTOIAWASE_YY, NSYDN_ORDER] last_record: ["00100002", "00100002", "2020", "00100002", "2020", 1]
3. 简化配置(推荐)
关闭use_raw_query_with_incremental,让Embulk自动生成多主键增量查询语句,避免手动编写SQL的错误:
# 移除use_raw_query_with_incremental配置 # use_raw_query_with_incremental: true table: TEST.TABLE1 select: [CTOIAWASE_NO, CTOIAWASE_YY, NSYDN_ORDER] incremental_columns: [CTOIAWASE_NO, CTOIAWASE_YY, NSYDN_ORDER] incremental: true last_record: ["00100002", "2020", 1]
4. 升级插件版本
当前使用的embulk-input-oracle版本为0.9.3,存在多主键增量的已知bug,建议升级到最新稳定版(如0.10.3),修复参数绑定和类型转换问题。
5. 验证最后一条记录同步
修复配置后,通过以下方式验证:
- 同步完成后检查目标数据是否包含最后一条记录
- 查看Embulk生成的SQL语句,确认增量条件正确
- 检查
last_record是否被正确更新为同步的最后一条记录的值
内容的提问来源于stack exchange,提问作者Saito Mieko

