【问题标题】:Join spark dataset using java使用 java 加入 spark 数据集
【发布时间】:2020-08-21 20:00:36
【问题描述】:

我有 2 个要合并的数据集:

数据集1(机器):

 String machineID:
 List<Integer> machineCat;(100,200,300)

数据集2(汽车):

 String carID:
 List<Integer> carCat;(30,200,100,300)

我基本上需要从 dataset1 获取 List machineCat 的每个项目,并检查它是否包含在 dataset2 的 List carCat 中。如果匹配,则将 2 个数据集组合如下:

最终数据集:

machineID,machineCat(100),carID,carCat(100)
machineID,machineCat(200),carID,carCat(200)
machineID,machineCat(300),carID,carCat(400)

关于如何在 java 中使用数据集连接的任何帮助。

查看带有 arrays_contain 的选项(如下所示)

  machine.foreachPartition((ForeachPartitionFunction<Machine>) iterator -> {

    while (iterator.hasNext()) {

        Machine machine = iterator.next();
        machine.getmachineCat().stream().filter(cat -> {

            LOG.info("matched");
            spark.sql(
                    "select * from machineDataset m"
                            + " join"
                            + " carDataset c "
                            + "where array_contains(m.machineCat,cat)");
            return true;
        });

    }
});

【问题讨论】:

  • 这太复杂了吗????
  • 你可以explode列出值,它会产生 machineId 和 machineCat (String,Integer) 列,对 dataset2 做同样的事情,然后通过 macineCat = carCat 加入,最后如果你想要 1 -元素数组使用array("machineCat")函数
  • @chlebek Dataset1(machine) 和 Dataset2(car) 包含的元素比列表和字符串多得多(超过 30 个)。您是否仍然认为爆炸列表会很容易。你能详细说明一下吗?
  • 这应该不会打扰,直到您爆炸超过一列

标签: java apache-spark dataset


【解决方案1】:
import static org.apache.spark.sql.functions.*; // before main class

Machine machine = new Machine("m1",Arrays.asList(100,200,300));
Car car = new Car("c1", Arrays.asList(30,200,100,300));

Dataset<Row> mDF= spark.createDataFrame(Arrays.asList(machine), Machine.class);
mDF.show();
Dataset<Row> cDF= spark.createDataFrame(Arrays.asList(car), Car.class);
cDF.show();

输出:

+---------------+---------+
|     machineCat|machineId|
+---------------+---------+
|[100, 200, 300]|       m1|
+---------------+---------+

+-------------------+-----+
|             carCat|catId|
+-------------------+-----+
|[30, 200, 100, 300]|   c1|
+-------------------+-----+

然后

Dataset<Row> mDF2 = mDF.select(col("machineId"),explode(col("machineCat")).as("machineCat"));
Dataset<Row> cDF2 = cDF.select(col("catId"),explode(col("carCat")).as("carCat"));
Dataset<Row> joinedDF = mDF2.join(cDF2).where(mDF2.col("machineCat").equalTo(cDF2.col("carCat")));
Dataset<Row> finalDF = joinedDF.select(col("machineId"),array(col("machineCat")), col("catId"),array(col("carCat")) );
finalDF.show();

最后:

+---------+-----------------+-----+-------------+
|machineId|array(machineCat)|catId|array(carCat)|
+---------+-----------------+-----+-------------+
|       m1|            [100]|   c1|        [100]|
|       m1|            [200]|   c1|        [200]|
|       m1|            [300]|   c1|        [300]|
+---------+-----------------+-----+-------------+

root
 |-- machineId: string (nullable = true)
 |-- array(machineCat): array (nullable = false)
 |    |-- element: integer (containsNull = true)
 |-- catId: string (nullable = true)
 |-- array(carCat): array (nullable = false)
 |    |-- element: integer (containsNull = true)

【讨论】:

  • 我的结果显示:[1172573,WrappedArray(141),549,WrappedArray(141)] [1172573,WrappedArray(1653),3155,WrappedArray(1653)] [1172573,WrappedArray(191), 1412,WrappedArray(191)] 如何展开这个 WrappedArray?
猜你喜欢
  • 2018-12-19
  • 1970-01-01
  • 2016-07-27
  • 2017-11-24
  • 2017-08-19
  • 2021-12-30
  • 1970-01-01
  • 1970-01-01
  • 2017-07-16
相关资源
最近更新 更多