【发布时间】: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