记csv、parquet数据预览一个bug的解决

2024-01-14 06:12

本文主要是介绍记csv、parquet数据预览一个bug的解决,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

文章目录

  • 一、概述
  • 二、实现过程
    • 1. 业务流程如图:
    • 2. 业务逻辑
    • 3. 运行结果
  • 三、bug现象
    • 1. 单元测试
    • 2.运行结果
  • 三、流程梳理
    • 1. 方向一
    • 2. 方向二

一、概述

工作中遇到通过sparksession解析csv、parquet文件并预览top100的需求。

二、实现过程

1. 业务流程如图:

hiveSQL读取数据
数据写入csv或parquet文件
预览csv或parquet文件top100数据

2. 业务逻辑

为了便于测试,我们下面以单元测试中模拟数据来说明


import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;import com.alibaba.fastjson.JSONObject;import lombok.extern.slf4j.Slf4j;@Slf4j
public class GroupingByDataTest
{static List<String> result = new ArrayList<>();@BeforeAllpublic static void init(){result.add("{\"student_no\":\"0204006\",\"student_name\":\"学生6\",\"field\":\"项目6\",\"value2\":\"6\",\"sex\":\"女\"}");result.add("{\"student_no\":\"0204006\",\"student_name\":\"学生6\",\"field\":\"项目6\",\"value2\":\"6\",\"sex\":\"女\"}");result.add("{\"student_no\":\"0204006\",\"student_name\":\"学生6\",\"field\":\"项目6\",\"value2\":\"6\",\"sex\":\"女\"}");result.add("{\"student_no\":\"0204006\",\"student_name\":\"学生6\",\"field\":\"项目6\",\"value2\":\"6\",\"sex\":\"女\"}");}@Testpublic void test002(){Map<Object, List<Object>> r = result.stream().map(s -> JSONObject.parseObject(s).entrySet()) // map.flatMap(m -> m.stream()) // flatMap.collect(Collectors.groupingBy(mp -> mp.getKey(), Collectors.mapping(x -> x.getValue(), Collectors.toList())));log.info("{}", r);}
}

3. 运行结果

 com.fly.lambda.GroupingByDataTest - {student_name=[学生6, 学生6, 学生6, 学生6], student_no=[0204006, 0204006, 0204006, 0204006], value2=[6, 6, 6, 6], field=[项目6, 项目6, 项目6, 项目6], sex=[女, 女, 女, 女]}

目前看来一切正常。

三、bug现象

实际测试过程中发现,hive数据仓库中的字段由于各种原因并不一定都有值,从而导致csv、parquet保存结果时字段为空

1. 单元测试


import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;import com.alibaba.fastjson.JSONObject;import lombok.extern.slf4j.Slf4j;@Slf4j
public class GroupingByDataTest
{static List<String> result = new ArrayList<>();@BeforeAllpublic static void init(){result.add("{\"student_name\":\"学生1\",\"student_no\":\"0204001\",\"field\":\"项目1\",                 \"sex\":\"男\"}");result.add("{\"student_name\":\"学生2\",\"student_no\":\"0204002\",\"field\":\"项目2\",\"value2\":\"2\"               }");result.add("{\"student_name\":\"学生3\",                           \"field\":\"项目3\",\"value2\":\"3\",\"sex\":\"女\"}");result.add("{                           \"student_no\":\"0204004\",\"field\":\"项目4\",\"value2\":\"4\",\"sex\":\"男\"}");result.add("{\"student_name\":\"学生5\",\"student_no\":\"0204005\",\"field\":\"项目5\",\"value2\":\"5\",\"sex\":\"女\"}");result.add("{\"student_no\":\"0204006\",\"student_name\":\"学生6\",\"field\":\"项目6\",\"value2\":\"6\",\"sex\":\"女\"}");}@Testpublic void test002(){Map<Object, List<Object>> r = result.stream().map(s -> JSONObject.parseObject(s).entrySet()) // map.flatMap(m -> m.stream()) // flatMap.collect(Collectors.groupingBy(mp -> mp.getKey(), Collectors.mapping(x -> x.getValue(), Collectors.toList())));log.info("{}", r);}
}

2.运行结果

 com.fly.lambda.GroupingByDataTest - {student_name=[学生1, 学生2, 学生3, 学生5, 学生6], student_no=[0204001, 0204002, 0204004, 0204005, 0204006], value2=[2, 3, 4, 5, 6], field=[项目1, 项目2, 项目3, 项目4, 项目5, 项目6], sex=[男, 女, 男, 女, 女]}

期望的结果为

 com.fly.lambda.GroupingByDataTest - before : {student_name=[学生1, 学生2, 学生3, null, 学生5, 学生6], student_no=[0204001, 0204002, null, 0204004, 0204005, 0204006], value2=[null, 2, 3, 4, 5, 6], field=[项目1, 项目2, 项目3, 项目4, 项目5, 项目6], sex=[男, null, 女, 男, 女, 女]}

三、流程梳理

解决这个问题有2个方向

1. 方向一

从数据来源解决,也就是 hiveSQL读取数据使用 coalsce 函数进行空值处理,实际去解决的过程中发现2个问题。

  1. 强制业务用户编辑hiveSQL时显式调用(用户体验太差,增加使用难度
  2. 不强制业务用户编辑hiveSQL时显式调用,后台接受到SQL后自动添加coalsce 函数(后台业务逻辑复杂,eg: 使用了条件语句、多表关联查询等等情况。不一而足,几乎没法妥善处理

2. 方向二

hiveSQL读取数据
数据写入csv或parquet文件
预览csv或parquet文件top100数据

hiveSQL读取数据、数据写入csv或parquet文件正常进行,不用特殊处理, 修改步骤3

分为2步骤,步骤1,遍历获取全部的key去重,步骤2,自动对缺失数据的key补充空值

核心代码如下:

@Testpublic void test003()throws IOException{// 取keysList<String> keys = result.stream().map(s -> JSONObject.parseObject(s).entrySet()).flatMap(m -> m.stream()).map(r -> r.getKey()).distinct().collect(Collectors.toList());keys.stream().forEach(log::info);Map<String, List<Object>> r = result.stream().map(s -> parse(s, keys)).flatMap(m -> m.stream()) // flatMap.collect(Collectors.groupingBy(mp -> mp.getKey(), Collectors.mapping(x -> x.getValue(), Collectors.toList())));log.info("before : {}", r);log.info("sorted : {}", new TreeMap<>(r));}/*** 设置value, 根据需要补充空值*/private Set<Entry<String, Object>> parse(String s, List<String> keys){JSONObject jsonObject = JSONObject.parseObject(s);keys.stream().forEach(key -> {if (!jsonObject.containsKey(key)){jsonObject.put(key, null);}});return jsonObject.entrySet();}

可以说,花比较小的成本,以比较少的代码变动,相对稳妥的解决了问题。


有任何问题和建议,都可以向我提问讨论,大家一起进步,谢谢!

-over-

这篇关于记csv、parquet数据预览一个bug的解决的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/604177

相关文章

大模型研发全揭秘:客服工单数据标注的完整攻略

在人工智能(AI)领域,数据标注是模型训练过程中至关重要的一步。无论你是新手还是有经验的从业者,掌握数据标注的技术细节和常见问题的解决方案都能为你的AI项目增添不少价值。在电信运营商的客服系统中,工单数据是客户问题和解决方案的重要记录。通过对这些工单数据进行有效标注,不仅能够帮助提升客服自动化系统的智能化水平,还能优化客户服务流程,提高客户满意度。本文将详细介绍如何在电信运营商客服工单的背景下进行

基于MySQL Binlog的Elasticsearch数据同步实践

一、为什么要做 随着马蜂窝的逐渐发展,我们的业务数据越来越多,单纯使用 MySQL 已经不能满足我们的数据查询需求,例如对于商品、订单等数据的多维度检索。 使用 Elasticsearch 存储业务数据可以很好的解决我们业务中的搜索需求。而数据进行异构存储后,随之而来的就是数据同步的问题。 二、现有方法及问题 对于数据同步,我们目前的解决方案是建立数据中间表。把需要检索的业务数据,统一放到一张M

关于数据埋点,你需要了解这些基本知识

产品汪每天都在和数据打交道,你知道数据来自哪里吗? 移动app端内的用户行为数据大多来自埋点,了解一些埋点知识,能和数据分析师、技术侃大山,参与到前期的数据采集,更重要是让最终的埋点数据能为我所用,否则可怜巴巴等上几个月是常有的事。   埋点类型 根据埋点方式,可以区分为: 手动埋点半自动埋点全自动埋点 秉承“任何事物都有两面性”的道理:自动程度高的,能解决通用统计,便于统一化管理,但个性化定

使用SecondaryNameNode恢复NameNode的数据

1)需求: NameNode进程挂了并且存储的数据也丢失了,如何恢复NameNode 此种方式恢复的数据可能存在小部分数据的丢失。 2)故障模拟 (1)kill -9 NameNode进程 [lytfly@hadoop102 current]$ kill -9 19886 (2)删除NameNode存储的数据(/opt/module/hadoop-3.1.4/data/tmp/dfs/na

异构存储(冷热数据分离)

异构存储主要解决不同的数据,存储在不同类型的硬盘中,达到最佳性能的问题。 异构存储Shell操作 (1)查看当前有哪些存储策略可以用 [lytfly@hadoop102 hadoop-3.1.4]$ hdfs storagepolicies -listPolicies (2)为指定路径(数据存储目录)设置指定的存储策略 hdfs storagepolicies -setStoragePo

Hadoop集群数据均衡之磁盘间数据均衡

生产环境,由于硬盘空间不足,往往需要增加一块硬盘。刚加载的硬盘没有数据时,可以执行磁盘数据均衡命令。(Hadoop3.x新特性) plan后面带的节点的名字必须是已经存在的,并且是需要均衡的节点。 如果节点不存在,会报如下错误: 如果节点只有一个硬盘的话,不会创建均衡计划: (1)生成均衡计划 hdfs diskbalancer -plan hadoop102 (2)执行均衡计划 hd

【Prometheus】PromQL向量匹配实现不同标签的向量数据进行运算

✨✨ 欢迎大家来到景天科技苑✨✨ 🎈🎈 养成好习惯,先赞后看哦~🎈🎈 🏆 作者简介:景天科技苑 🏆《头衔》:大厂架构师,华为云开发者社区专家博主,阿里云开发者社区专家博主,CSDN全栈领域优质创作者,掘金优秀博主,51CTO博客专家等。 🏆《博客》:Python全栈,前后端开发,小程序开发,人工智能,js逆向,App逆向,网络系统安全,数据分析,Django,fastapi

如何解决线上平台抽佣高 线下门店客流少的痛点!

目前,许多传统零售店铺正遭遇客源下降的难题。尽管广告推广能带来一定的客流,但其费用昂贵。鉴于此,众多零售商纷纷选择加入像美团、饿了么和抖音这样的大型在线平台,但这些平台的高佣金率导致了利润的大幅缩水。在这样的市场环境下,商家之间的合作网络逐渐成为一种有效的解决方案,通过资源和客户基础的共享,实现共同的利益增长。 以最近在上海兴起的一个跨行业合作平台为例,该平台融合了环保消费积分系统,在短

烟火目标检测数据集 7800张 烟火检测 带标注 voc yolo

一个包含7800张带标注图像的数据集,专门用于烟火目标检测,是一个非常有价值的资源,尤其对于那些致力于公共安全、事件管理和烟花表演监控等领域的人士而言。下面是对此数据集的一个详细介绍: 数据集名称:烟火目标检测数据集 数据集规模: 图片数量:7800张类别:主要包含烟火类目标,可能还包括其他相关类别,如烟火发射装置、背景等。格式:图像文件通常为JPEG或PNG格式;标注文件可能为X

pandas数据过滤

Pandas 数据过滤方法 Pandas 提供了多种方法来过滤数据,可以根据不同的条件进行筛选。以下是一些常见的 Pandas 数据过滤方法,结合实例进行讲解,希望能帮你快速理解。 1. 基于条件筛选行 可以使用布尔索引来根据条件过滤行。 import pandas as pd# 创建示例数据data = {'Name': ['Alice', 'Bob', 'Charlie', 'Dav