Java中Spark遍历DataFrame分区生成HBase Delete对象取值报错如何解决
问题原因
- 多余的Row包装逻辑:
iterator.next()本身返回的就是Row实例,额外用RowFactory.create包装会将原Row作为新Row的唯一元素,无法直接读取原数据的hbase_key字段 - 字段取值错误:直接对Row对象调用
String.valueOf()会得到Row对象的toString结果,而非hbase_key字段的值 - 分布式变量安全问题:如果
deleteList是定义在Driver端的集合,直接在foreachPartition中修改不会同步回Driver,因为算子逻辑在Executor端运行时拿到的是变量的序列化副本
正确实现
首先确认hbase_key字段的类型,以最常见的String类型为例,正确代码如下:
df.foreachPartition((ForeachPartitionFunction<Row>) iterator -> { // HBase连接、deleteList都需要在分区内部初始化,避免序列化问题和线程安全问题 List<Delete> deleteList = new ArrayList<>(); // 提前获取hbase_key字段索引,避免循环内重复查找损耗性能 int keyIndex = -1; while (iterator.hasNext()) { Row row = iterator.next(); if (keyIndex == -1) { keyIndex = row.fieldIndex("hbase_key"); } // 按索引读取hbase_key的值,可根据实际字段类型调整取值方法 String hbaseKey = row.getString(keyIndex); deleteList.add(new Delete(Bytes.toBytes(hbaseKey))); } // 此处直接实现批量提交HBase删除请求的逻辑即可,不要尝试把deleteList传回Driver端 });
类型适配说明
如果hbase_key不是String类型,替换对应的Row取值方法即可:
- 长整型:
row.getLong(keyIndex) - 字节数组:
row.getBinary(keyIndex) - 其他类型可直接对应Spark Row的API调整
内容的提问来源于stack exchange,提问作者Asynchronous
相关产品推荐
相关产品推荐

