Spark 1.6中registerTempTable()未注册全部记录的问题求助
嘿,针对你在Spark 1.6里遇到的这个问题——从RDB读数据注册临时表后,生成两个Parquet文件(间隔两小时)且查询临时表拿不到全量记录,我帮你拆解原因并给出针对性的解决方案:
核心原因拆解
Spark延迟计算的坑
Spark 1.6的DataFrame操作是延迟计算的,只有遇到count()、collect()这类action操作才会真正触发从RDB拉取数据。如果你的代码只注册了临时表就直接写Hive,可能写操作只触发了部分数据的拉取,剩下的分区因为资源不足、RDB端性能瓶颈(比如锁、慢查询)延迟加载,导致两小时后才生成第二个Parquet文件。而你查询临时表时,可能只读取了已经完成加载的部分数据,自然拿不全。临时表的生命周期与数据一致性问题
临时表是和当前Spark Session绑定的,它的数据依赖于背后的DataFrame计算状态。如果写Hive的Job还在异步执行(虽然Spark默认是同步,但极端情况下可能因为资源调度出现延迟),此时查询临时表,只能拿到已经完成计算的那部分数据。RDB读取的并行度不合理
如果从RDB读数据时没有设置合理的分区,单分区数据量过大,会导致该分区的读取和写入耗时极长,拖慢整体进度,出现分批生成Parquet文件的情况。
针对性解决方案
1. 强制触发全量数据加载后再注册临时表
在注册临时表前,先调用action操作把全量数据从RDB拉到Spark中,确保临时表的数据是完整的:
// 从RDB读取数据 val jdbcDF = sqlContext.read.jdbc( "jdbc:mysql://your-db-host:3306/dbname", "employees", connectionProps ) // 触发全量数据加载(count()是action操作,会阻塞直到数据拉取完成) jdbcDF.count() // 再注册临时表 jdbcDF.registerTempTable("temp_employees")
2. 确保写Hive操作完全执行后再查询
Spark的INSERT INTO是action操作,但如果你的代码没有等待它执行完毕就去查询,可能会拿到不完整的数据。可以用collect()强制阻塞等待Job完成:
// 执行插入Hive表的操作并等待完成 sqlContext.sql("INSERT INTO hive_employees SELECT * FROM temp_employees").collect() // 之后再查询临时表或Hive表 val fullData = sqlContext.sql("SELECT * FROM temp_employees").collect()
3. 优化RDB读取的并行度
用分区参数拆分RDB读取任务,避免单分区数据过大:
val connectionProps = new Properties() connectionProps.put("user", "your-username") connectionProps.put("password", "your-password") // 按employee_id分区,分成10个并行任务读取 val jdbcDF = sqlContext.read.jdbc( "jdbc:mysql://your-db-host:3306/dbname", "employees", "employee_id", // 分区字段(需是数值型) 1, // 分区下界 10000, // 分区上界 10, // 分区数 connectionProps )
这样能让数据读取更均衡,避免单个任务耗时过长导致的延迟写入。
4. 直接写入Hive表(替代临时表方案)
虽然Spark 1.6和Hive的兼容性有限,但可以尝试直接用saveAsTable写入,避开临时表的生命周期问题:
// 直接将DataFrame写入Hive表(append模式) jdbcDF.write.mode("append").format("parquet").saveAsTable("hive_employees")
如果你的Hive表是预创建的,确保存储格式是Parquet,这样兼容性更好。
总结
本质问题是Spark延迟计算导致数据加载和写入不是一次性完成,通过强制全量加载、等待Job执行完毕、优化读取并行度这几点,基本能解决Parquet分批生成和查询不全的问题。
内容的提问来源于stack exchange,提问作者regiea

