发布时间:2026/9/5 15:15:57
Flink SQL 基于Update流出现空值无法过滤问题 问题背景问题描述基于Flink-CDC(2.3.0 Flink 1.15) Flink SQL的实时计算作业在运行一段时间后突然发现插入数据库的计算结果发生部分主键属性发生失败导致后续计算结果无法插入 超过失败次数失败的情况问题报错Caused by: java.sql.BatchUpdateException: Batch entry 0 INSERT INTO dm_hljy.dws_table_name (op_date, school_year, campus_name, school_name, depart_name, total_opfare, ids, update_time) VALUES (2024-03-11 00:00:0008, 2023, xxxx, xxxx学校, xxxx小学部, 203333300000, 57, 2024-03-21 09:31:08.4708) ON DUPLICATE KEY UPDATE school_yearVALUES(school_year), total_opfareVALUES(total_opfare), idsVALUES(ids), update_timeVALUES(update_time) was aborted: ERROR: dn_6007_6008: null value in column depart_name violates not-null constraint Call getNextException to see other errors in the batch. at com.huawei.gauss200.jdbc.jdbc.BatchResultHandler.handleCompletion(BatchResultHandler.java:171) ~[huaweicloud-dws-jdbc-8.1.1.1-200.jar:?] at com.huawei.gauss200.jdbc.core.v3.QueryExecutorImpl.executeBatch(QueryExecutorImpl.java:586) ~[huaweicloud-dws-jdbc-8.1.1.1-200.jar:?] at com.huawei.gauss200.jdbc.jdbc.PgStatement.executeBatch(PgStatement.java:883) ~[huaweicloud-dws-jdbc-8.1.1.1-200.jar:?] at com.huawei.gauss200.jdbc.jdbc.PgPreparedStatement.executeBatch(PgPreparedStatement.java:1580) ~[huaweicloud-dws-jdbc-8.1.1.1-200.jar:?] at org.apache.flink.connector.jdbc.statement.FieldNamedPreparedStatementImpl.executeBatch(FieldNamedPreparedStatementImpl.java:65) ~[flink-connector-jdbc-1.15.0-h0.cbu.mrs.320.r33.jar:1.15.0-h0.cbu.mrs.320.r33] at org.apache.flink.connector.jdbc.internal.executor.TableSimpleStatementExecutor.executeBatch(TableSimpleStatementExecutor.java:64) ~[flink-connector-jdbc-1.15.0-h0.cbu.mrs.320.r33.jar:1.15.0-h0.cbu.mrs.320.r33] at org.apache.flink.connector.jdbc.internal.executor.TableBufferReducedStatementExecutor.executeBatch(TableBufferReducedStatementExecutor.java:101) ~[flink-connector-jdbc-1.15.0-h0.cbu.mrs.320.r33.jar:1.15.0-h0.cbu.mrs.320.r33] at org.apache.flink.connector.jdbc.internal.JdbcOutputFormat.attemptFlush(JdbcOutputFormat.java:266) ~[flink-connector-jdbc-1.15.0-h0.cbu.mrs.320.r33.jar:1.15.0-h0.cbu.mrs.320.r33] at org.apache.flink.connector.jdbc.internal.JdbcOutputFormat.flush(JdbcOutputFormat.java:236) ~[flink-connector-jdbc-1.15.0-h0.cbu.mrs.320.r33.jar:1.15.0-h0.cbu.mrs.320.r33] at org.apache.flink.connector.jdbc.internal.JdbcOutputFormat.lambda$open$0(JdbcOutputFormat.java:159) ~[flink-connector-jdbc-1.15.0-h0.cbu.mrs.320.r33.jar:1.15.0-h0.cbu.mrs.320.r33] at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) ~[?:1.8.0_332] at java.util.concurrent.FutureTask.runAndReset(FutureTask.java:308) ~[?:1.8.0_332] at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$301(ScheduledThreadPoolExecutor.java:180) ~[?:1.8.0_332] at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:294) ~[?:1.8.0_332] at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) ~[?:1.8.0_332] at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) ~[?:1.8.0_332] ... 1 more Caused by: com.huawei.gauss200.jdbc.util.PSQLException: ERROR: dn_6007_6008: null value in column depart_name violates not-null constraint at com.huawei.gauss200.jdbc.core.v3.QueryExecutorImpl.receiveErrorResponse(QueryExecutorImpl.java:2856) ~[huaweicloud-dws-jdbc-8.1.1.1-200.jar:?] at com.huawei.gauss200.jdbc.core.v3.QueryExecutorImpl.processResults(QueryExecutorImpl.java:2587) ~[huaweicloud-dws-jdbc-8.1.1.1-200.jar:?] at com.huawei.gauss200.jdbc.core.v3.QueryExecutorImpl.executeBatch(QueryExecutorImpl.java:575) ~[huaweicloud-dws-jdbc-8.1.1.1-200.jar:?] at com.huawei.gauss200.jdbc.jdbc.PgStatement.executeBatch(PgStatement.java:883) ~[huaweicloud-dws-jdbc-8.1.1.1-200.jar:?] at com.huawei.gauss200.jdbc.jdbc.PgPreparedStatement.executeBatch(PgPreparedStatement.java:1580) ~[huaweicloud-dws-jdbc-8.1.1.1-200.jar:?] at org.apache.flink.connector.jdbc.statement.FieldNamedPreparedStatementImpl.executeBatch(FieldNamedPreparedStatementImpl.java:65) ~[flink-connector-jdbc-1.15.0-h0.cbu.mrs.320.r33.jar:1.15.0-h0.cbu.mrs.320.r33] at org.apache.flink.connector.jdbc.internal.executor.TableSimpleStatementExecutor.executeBatch(TableSimpleStatementExecutor.java:64) ~[flink-connector-jdbc-1.15.0-h0.cbu.mrs.320.r33.jar:1.15.0-h0.cbu.mrs.320.r33] at org.apache.flink.connector.jdbc.internal.executor.TableBufferReducedStatementExecutor.executeBatch(TableBufferReducedStatementExecutor.java:101) ~[flink-connector-jdbc-1.15.0-h0.cbu.mrs.320.r33.jar:1.15.0-h0.cbu.mrs.320.r33] at org.apache.flink.connector.jdbc.internal.JdbcOutputFormat.attemptFlush(JdbcOutputFormat.java:266) ~[flink-connector-jdbc-1.15.0-h0.cbu.mrs.320.r33.jar:1.15.0-h0.cbu.mrs.320.r33] at org.apache.flink.connector.jdbc.internal.JdbcOutputFormat.flush(JdbcOutputFormat.java:236) ~[flink-connector-jdbc-1.15.0-h0.cbu.mrs.320.r33.jar:1.15.0-h0.cbu.mrs.320.r33] at org.apache.flink.connector.jdbc.internal.JdbcOutputFormat.lambda$open$0(JdbcOutputFormat.java:159) ~[flink-connector-jdbc-1.15.0-h0.cbu.mrs.320.r33.jar:1.15.0-h0.cbu.mrs.320.r33] at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) ~[?:1.8.0_332] at java.util.concurrent.FutureTask.runAndReset(FutureTask.java:308) ~[?:1.8.0_332] at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$301(ScheduledThreadPoolExecutor.java:180) ~[?:1.8.0_332] at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:294) ~[?:1.8.0_332] at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) ~[?:1.8.0_332] at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) ~[?:1.8.0_332定位定位思路1.方向一怀疑数据库插入存在数据处理时造成数据处理出现空值的情况即数据本身不为空但是数据插入却出现了空 2.方向二Flink-SQL在消费kafka数据时存在了空值故加工的数据计算结果存在空值定位过程因插入数据库定位比较麻烦且数据库已经设置该字段为主属性故出现插入时处理为空值的概率较小。故先从较为简单的Flink SQL查询数据定位方法一查询该字段为空的记录,待作业执行完成后未查询到空值对应记录select select * from table_name where depart_name is null or depart_name or char_length(depart_name) 0;因考虑到使用Flink-CDC进行变更数据捕获故对应的update流存在-U,U,-D,I记录因此随着插入记录存在空值被记录进去的情况故采用view的方式先将宽表的加工、关联方式创建为view然后进行空值的过滤。实施如下create view view_prd as select a.* ,b.* from a join b on a.id b.id select * from view_prd where depart_name is null or depart_name or char_length(depart_name) 0;通过查询结果发现存在最后一条记录存在空值的原因往源头定位发现该字段之前为空后面进行更新填充到值出现-U记录导致数据插入持续失败原因因为flink-SQL消费的数据时kafka topicflink以upsert-kafka形式的connector进行写入故存在changelog 流中数据更新存在-UU的记录按照Key进行区分唯一条记录value 为空(-U)的记录kafka也导致出现空值解决通过在DWS宽表创建一层View如上)在写入DWS宽表的kafka topic之前现将该字段空值过滤即可排除空值涉及记录被纳入结果指标计算的范围中

相关新闻

2026/9/5 15:10:57

AI绘画角色一致性实战:LoRA与ControlNet打造多角度骑士西奥

最近在AI绘画圈子里,一个有趣的现象正在发生:许多创作者不再满足于生成单张精美的图片,而是开始热衷于围绕同一个角色,从不同角度、不同情境进行“角色深度塑造”。这背后反映的,其实是AI绘画技术从“炫技”走向“实用…

2026/9/5 15:10:57

毕业论文修改的进阶指南:从查重达标到质量跃升

在撰写毕业论文的过程中,文本修改是一个不可避免的环节。面对不同的修改方式,我常常感到困惑:是使用传统的同义词替换,还是借助通用大模型辅助改写,或者使用专门的论文文本处理工具?各自适合怎样的任务&…

2026/9/5 15:56:01

VLDB 2026:微软研究院两篇论文获认可,数据库技术的未来风向标

如果你一直在关注数据库和系统领域的技术趋势,最近微软研究院传来的消息值得停下来看一眼:两篇论文获得了 VLDB 2026 的认可。对非学术圈的开发者来说,“论文获奖”听起来像是离日常开发很远的事。但 VLDB 不是普通会议——它是数据库与数据管…

2026/9/5 15:56:01

FLUENT17.0流体仿真工程实践:从参数物理意义到工业级收敛

简介:本资源是《FLUENT17.0流体仿真从入门到精通》配套实践文件包,面向CFD初学者及工程仿真从业人员,系统解决流体建模、网格划分、物理模型设置、求解计算与后处理分析等核心学习难点。压缩包共253个文件,涵盖39个cas&#xff08…

2026/9/5 15:56:01

.NET 8全栈开源在线考试系统:跨平台、多数据库与高并发架构实践

简介:星期八在线考试系统是一套面向高校教务部门、职业院校及企业培训中心的开源免费企业级教学管理解决方案,聚焦在线考试全流程数字化,解决题库建设、智能组卷、防作弊监考、自动阅卷与多维成绩分析等核心痛点。资源包共2000个文件&#xf…

2026/9/5 2:46:54

vSound小提琴数字处理器实操指南:从接线到演出的完整配置

电小提琴或者原声小提琴插电演出,第一个绕不开的坎就是声音难听。原声琴的共鸣和空气感一旦进了拾音器,出来的往往是一坨干瘪、发尖、带着奇怪塑料味的信号。我当初第一次把琴接上乐队调音台,直接被主唱吐槽"你这声音像在锯钢丝"。…

2026/9/5 2:46:52

传感器接口IC如何攻克生物化学传感的微弱信号难题?

1. 从电极到比特流:为什么生物化学传感必须依赖专用接口IC 做生物化学传感的人都有过类似的经历:明明传感器本身性能很好,信号输出却一塌糊涂——噪声大、漂移明显、重复性差,怎么调都达不到预期。很多时候问题并不在传感器&#…

2026/9/5 2:44:34

STM32F411CEU6多通道ADC采集:扫描模式+DMA实现详解

1. 多通道 ADC 的用武之地把“Multichannel ADC”和“STM32F411CEU6”这两个关键字放在一起,其实就是嵌入式开发里最常遇到的一类需求:用一块不算贵的 MCU,同时采集多路模拟信号。STM32F411CEU6 是 48 引脚的 Cortex-M4F 主控,主频…

2026/9/5 0:04:47

流式背压机制:避免前端渲染卡死与内存暴涨的滑动窗口限流

流式背压机制:避免前端渲染卡死与内存暴涨的滑动窗口限流在大模型流式输出(Streaming)与智能体实时推流的架构中,生产环境中经常出现一种“上下游生产消费速率严重失衡”的极端情况: 生产端极速产出:大模型…

2026/9/5 2:45:13

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/5 2:30:42

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/5 2:46:50

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…