Scala读取超4096字符IBM MQ消息出现损坏记录如何解决
问题根因
你遇到的消息被截断为4095字符的问题,根因来自四个层面:
- IBM MQ 9.2默认的通道、队列、队列管理器的
MAXMSGL参数为4096字节,超过该长度的消息会被默认截断 - 你使用的
spark-jms-receiver系列依赖版本默认的消息接收缓冲区硬编码为4KB,无法接收超过长度的消息 - 依赖包存在版本冲突:你引入了适配Spark 1.5.2的
spark-core_2.11-1.5.2.logging.jar,和当前运行的Spark 2.4.5核心逻辑不兼容,会干扰消息读取逻辑 - 你传入的自定义
converter转换器可能存在硬编码的字节数组长度限制,仅读取前4096字节的消息内容
修复方案
1. 调整IBM MQ侧配置
将MQ队列管理器、接收队列、对应连接通道的MAXMSGL参数统一调整为大于你最大消息长度的值,建议设置为1048576(1MB)即可覆盖绝大多数场景,同时在MQ连接工厂中添加缓冲区配置:
// 在MQConsumerFactory初始化逻辑中添加以下配置 connectionFactory.setReceiveBufferSize(1048576); connectionFactory.setClientReconnectOptions(WMQConstants.WMQ_CLIENT_RECONNECT_YES);
2. 修复依赖冲突与JMS Receiver配置
- 首先从
JARLIST中删除冗余的低版本依赖:spark-core_2.11-1.5.2.logging.jar,同时如果spark-jms-receiver-0.1.2-s_2.11.jar和spark-mq-jms-receiver_2.11-0.0.1-SNAPSHOT.jar功能重复,仅保留适配Spark 2.4.5版本的单个依赖即可 - 若你使用的JMS Receiver支持配置最大消息长度,在调用
createAsynchronousJmsQueueStream时传入maxMessageLength=1048576参数,覆盖默认4KB的缓冲区限制 - 检查自定义
converter逻辑,删除硬编码的4096字节数组限制,改为全量读取JMS消息字节流:// 示例正确的消息转换逻辑 def jmsMessageConverter(message: Message): String = { message match { case textMsg: TextMessage => textMsg.getText case bytesMsg: BytesMessage => val baos = new ByteArrayOutputStream() val buffer = new Array[Byte](1024) var len = 0 while ({len = bytesMsg.read(buffer); len != -1}) { baos.write(buffer, 0, len) } baos.toString(StandardCharsets.UTF_8.name()) } }
3. Spark端配置调整
Spark本身没有默认的字符串长度限制,你可以在spark-submit命令中添加以下配置避免JSON解析、序列化阶段出现长度限制:
# 添加到spark-submit参数中 --conf spark.driver.maxResultSize=0 \ --conf spark.sql.json.maxStringLength=10485760 \ --conf spark.serializer=org.apache.spark.serializer.KryoSerializer
其中spark.sql.json.maxStringLength控制Spark SQL解析JSON时允许的最大字符串长度,默认值为1MB,调整为10MB即可满足需求;spark.driver.maxResultSize=0表示取消driver端结果收集的大小限制。
内容的提问来源于stack exchange,提问作者Chia
相关产品推荐
相关产品推荐

