【发布时间】:2017-09-22 20:42:25
【问题描述】:
我有一组非常烦人的文件结构如下:
userId string,
eventType string,
source string,
errorCode string,
startDate timestamp,
endDate timestamp
每个文件的每个 eventId 可以包含任意数量的记录,具有不同的 eventTypes 和来源,以及每个文件的不同代码和开始/结束日期。
在 Hive 或 Spark 中是否有一种方法可以在 userId 上将所有这些组合在一起,有点像键值,其中值是与 userId 关联的所有字段的列表?具体来说,我希望它由 eventType 和 source 键入。基本上我想用表格长度换取宽度,有点像数据透视表。我的目标是最终以 Apache Parquet 或 Avro 文件格式存储,以便将来进行更快速的分析。
这是一个例子:
来源数据:
userId, eventType, source, errorCode, startDate, endDate
552113, 'ACK', 'PROVIDER', 0, '2017-09-01 12:01:45.432', '2017-09-01 12:01:45.452'
284723, 'ACK', 'PROVIDER', 0, '2017-09-01 12:01:45.675', '2017-09-01 12:01:45.775'
552113, 'TRADE', 'MERCH', 0, '2017-09-01 12:01:47.221', '2017-09-01 12:01:46.229'
552113, 'CHARGE', 'MERCH', 0, '2017-09-01 12:01:48.123', '2017-09-01 12:01:48.976'
284723, 'REFUND', 'MERCH', 1, '2017-09-01 12:01:48.275', '2017-09-01 12:01:48.947'
552113, 'CLOSE', 'PROVIDER', 0, '2017-09-01 12:01:49.908', '2017-09-01 12:01:50.623'
284723, 'CLOSE', 'PROVIDER', 0, '2017-09-01 12:01:50.112', '2017-09-01 12:01:50.777'
目标:
userId, eventTypeAckProvider, sourceAckProvider, errorCodeAckProvider, startDateAckProvider, endDateAckProvider, eventTypeTradeMerch, sourceTradeMerch, errorCodeTradeMerch, startDateTradeMerch, endDateTradeMerch, eventTypeChargeMerch, sourceChargeMerch, errorCodeChargeMerch, startDateChargeMerch, endDateChargeMerch, eventTypeCloseProvider, sourceCloseProvider, errorCodeCloseProvider, startDateCloseProvider, endDateCloseProvider, eventTypeRefundMerch, sourceRefundMerch, errorCodeRefundMerch, startDateRefundMerch, endDateRefundMerch
552113, 'ACK', 'PROVIDER', 0, '2017-09-01 12:01:45.432', '2017-09-01 12:01:45.452', 'TRADE', 'MERCH', 0, '2017-09-01 12:01:47.221', '2017-09-01 12:01:46.229', 'CHARGE', 'MERCH', 0, '2017-09-01 12:01:48.123', '2017-09-01 12:01:48.976', 'CLOSE', 'PROVIDER', 0, '2017-09-01 12:01:49.908', '2017-09-01 12:01:50.623', NULL, NULL, NULL, NULL, NULL
284723, 'ACK', 'PROVIDER', 0, '2017-09-01 12:01:45.675', '2017-09-01 12:01:45.775', NULL, NULL, NULL, NULL, NULL, NULL, NULL, NULL, NULL, NULL, 'CLOSE', 'PROVIDER', 0, '2017-09-01 12:01:50.112', '2017-09-01 12:01:50.777', 'REFUND', 'MERCH', 1, '2017-09-01 12:01:48.275', '2017-09-01 12:01:48.947'
字段名或顺序无所谓,只要我能区分就行。
我已经尝试了两种方法来让它工作:
- 从表中手动选择每个组合并加入主数据集。这工作得很好,并行化也很好,但不允许关键字段有任意数量的值,并且需要预定义架构。
- 使用 Spark 创建键值字典,其中每个值都是字典。基本上遍历数据集,如果字典不存在,则向字典添加一个新键,如果它不存在,则为该条目向值字典添加一个新字段。这工作得很好,但是非常慢并且不能很好地并行化,如果可以的话。我也不确定这是否是 Avro/Parquet 兼容格式。
这两种方法有什么替代方法吗?甚至比我的目标更好的结构?
【问题讨论】:
标签: hive pyspark avro emr parquet