【问题标题】:Save Spark Dataframe into Elasticsearch - Can’t handle type exception将 Spark Dataframe 保存到 Elasticsearch - 无法处理类型异常
【发布时间】:2015-12-16 11:47:35
【问题描述】:

我设计了一个简单的作业来从 MySQL 读取数据并使用 Spark 将其保存在 Elasticsearch 中。

代码如下:

JavaSparkContext sc = new JavaSparkContext(
        new SparkConf().setAppName("MySQLtoEs")
                .set("es.index.auto.create", "true")
                .set("es.nodes", "127.0.0.1:9200")
                .set("es.mapping.id", "id")
                .set("spark.serializer", KryoSerializer.class.getName()));

SQLContext sqlContext = new SQLContext(sc);

// Data source options
Map<String, String> options = new HashMap<>();
options.put("driver", MYSQL_DRIVER);
options.put("url", MYSQL_CONNECTION_URL);
options.put("dbtable", "OFFERS");
options.put("partitionColumn", "id");
options.put("lowerBound", "10001");
options.put("upperBound", "499999");
options.put("numPartitions", "10");

// Load MySQL query result as DataFrame
LOGGER.info("Loading DataFrame");
DataFrame jdbcDF = sqlContext.load("jdbc", options);
DataFrame df = jdbcDF.select("id", "title", "description",
        "merchantId", "price", "keywords", "brandId", "categoryId");
df.show();
LOGGER.info("df.count : " + df.count());
EsSparkSQL.saveToEs(df, "offers/product");

您可以看到代码非常简单。它将数据读入 DataFrame,选择一些列,然后执行 count 作为对 Dataframe 的基本操作。到目前为止一切正常。

然后它尝试将数据保存到 Elasticsearch 中,但由于无法处理某些类型而失败。可以看到错误日志here

我不确定为什么它不能处理那种类型。 有人知道为什么会这样吗?

我正在使用 Apache Spark 1.5.0、Elasticsearch 1.4.4 和 elaticsearch-hadoop 2.1.1

编辑:

  • 我已使用示例数据集和源代码更新了要点链接。
  • 我还尝试使用邮件列表中@costin 提到的elasticsearch-hadoop dev builds

【问题讨论】:

    标签: elasticsearch apache-spark elasticsearch-hadoop apache-spark-1.5


    【解决方案1】:

    这个问题的答案很棘手,但感谢samklr,我设法弄清楚了问题所在。

    尽管如此,解决方案并不简单,可能会考虑一些“不必要的”转换。

    首先让我们谈谈序列化

    在Spark数据序列化和函数序列化中,序列化有两个方面需要考虑。在这种情况下,它是关于数据序列化和反序列化的。

    从 Spark 的角度来看,唯一需要做的就是设置序列化 - Spark 默认依赖 Java 序列化,这很方便,但效率相当低。这就是Hadoop本身引入了自己的序列化机制和自己的类型——即Writables的原因。因此,InputFormatOutputFormats 必须返回 Writables,Spark 开箱即用不理解。

    使用 elasticsearch-spark 连接器,必须启用一种不同的序列化 (Kryo),它会自动处理转换并且非常高效。

    conf.set("spark.serializer","org.apache.spark.serializer.KryoSerializer")
    

    即使 Kryo 不要求类实现要序列化的特定接口,这意味着 POJO 可以在 RDD 中使用,除了启用 Kryo 序列化之外,无需任何进一步的工作。

    也就是说,@samklr 向我指出 Kryo 需要在使用它们之前注册类。

    这是因为 Kryo 写入了对正在序列化的对象的类的引用(为每个写入的对象写入一个引用),如果该类已注册,则它只是一个整数标识符,否则是完整的类名。 Spark 代表您注册 Scala 类和许多其他框架类(如 Avro Generic 或 Thrift 类)。

    向 Kryo 注册课程非常简单。创建 KryoRegistrator 的子类,并重写 registerClasses() 方法:

    public class MyKryoRegistrator implements KryoRegistrator, Serializable {
        @Override
        public void registerClasses(Kryo kryo) {
            // Product POJO associated to a product Row from the DataFrame            
            kryo.register(Product.class); 
        }
    }
    

    最后,在您的驱动程序中,将 spark.kryo.registrator 属性设置为您的 KryoRegistrator 实现的完全限定类名:

    conf.set("spark.kryo.registrator", "MyKryoRegistrator")
    

    其次,即使设置了 Kryo 序列化程序并注册了类,对 Spark 1.5 进行了更改,但由于某种原因,Elasticsearch 无法反序列化 Dataframe,因为它无法推断将 Dataframe 的 SchemaType 插入连接器。

    所以我不得不将 Dataframe 转换为 JavaRDD

    JavaRDD<Product> products = df.javaRDD().map(new Function<Row, Product>() {
        public Product call(Row row) throws Exception {
            long id = row.getLong(0);
            String title = row.getString(1);
            String description = row.getString(2);
            int merchantId = row.getInt(3);
            double price = row.getDecimal(4).doubleValue();
            String keywords = row.getString(5);
            long brandId = row.getLong(6);
            int categoryId = row.getInt(7);
            return new Product(id, title, description, merchantId, price, keywords, brandId, categoryId);
        }
    });
    

    现在数据已经准备好写入 elasticsearch 了:

    JavaEsSpark.saveToEs(products, "test/test");
    

    参考资料:

    • Elasticsearch 的 Apache Spark 支持 documentation
    • Hadoop 权威指南,第 19 章。Spark,编辑。 4 – 汤姆·怀特。
    • 用户samklr

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2017-01-26
      • 1970-01-01
      • 2019-03-29
      • 2016-05-26
      • 2020-03-08
      • 2018-03-02
      • 1970-01-01
      相关资源
      最近更新 更多