【问题标题】:Spark Cassandra connector with Java for read用于读取的带有 Java 的 Spark Cassandra 连接器
【发布时间】:2020-10-31 08:50:20
【问题描述】:

要求:- 我在 cassandra 中保存了数据,并且每小时我需要根据记录的更新来计算一些分数。我通过在数据集上使用 show() 方法看到数据正确

读取数据的代码如下:-

Dataset<DealFeedSchema> dealFeedSchemaDataset = session.read()
     .format(Constants.SPARK_CASSANDRA_SOURCE_PATH)
     .option(Constants.KEY_SPACE, Constants.CASSANDRA_KEY_SPACE)
     .option(Constants.TABLE, Constants.CASSANDRA_DEAL_TABLE_SPACE)
     .option(Constants.DATE_FORMAT, "yyyy-MM-dd HH:mm:ss")
     .schema(DealFeedSchema.getDealFeedSchema())
     .load()
     .as(Encoders.bean(DealFeedSchema.class));
dealFeedSchemaDataset.show();

show 的输出如下:

+-------+----------+-------------+--------------------+-----------+------------+----------+------------------------+---------------+-----------+-------------------+-------------------+----------+------------+-------------------+----------------+-------------------+-------------+----------+--------------------+----------+-------------------------+---------------+----------------+---------------+--------------+--------------+-----+
|deal_id| deal_name|deal_category|           deal_tags|growth_tags|deal_tag_ids|deal_price|deal_discount_percentage|deal_group_size|deal_active|    deal_start_time|        deal_expiry|product_id|product_name|product_description|product_category|product_category_id|product_price|hero_image|      product_images| video_url|video_thumbnail_image_url|deal_like_count|deal_share_count|deal_view_count|deal_buy_count|weighted_score|boost|
+-------+----------+-------------+--------------------+-----------+------------+----------+------------------------+---------------+-----------+-------------------+-------------------+----------+------------+-------------------+----------------+-------------------+-------------+----------+--------------------+----------+-------------------------+---------------+----------------+---------------+--------------+--------------+-----+
|      4|7h12349961|          mqw|[under999, under3...|         []|          []|    4969.0|                    null|       95166551|          1|2020-07-08 14:48:57|2020-07-18 14:48:57|4725457233|  kao62ggnm7|         32h64e356z|      jnnh29zr1f|               null|       6651.0|86kk7s34yr|[dSt4P79, i4WXOHb...|d6tag27924|               4j1l36lp17|           null|            null|           null|          null|          null| null|

所以当我在 dealFeedSchemaDataset 上使用 map/foreach 时发生了奇怪的事情,数据似乎不正确我将 deal_start_time 的列值作为当前系统时间,如下所示,不知道这是如何改变的。

即使在下面的行也给出了同样的问题:

dealFeedSchemaDataset.select(
      functions.col("deal_start_time")).as(Encoders.bean(DateTime.class))
.collectAsList().forEach(schema -> System.out.println(schema));
2020-07-10T20:21:47.895+05:30

有人可以帮我解决我做错了什么吗?

【问题讨论】:

  • 如果您在加载数据时已经这样做了,为什么还需要再次执行as。此外,问题可能来自在.as 中使用了不正确的类型.select -> 该字段很可能具有timestamp 类型,您需要使用java.sql.Timestampjava.sql.Date。 .
  • 嗨@AlexOtt 是的,你是对的,我正在对代码进行反复试验。 java.sql.timestamp 是我一直在寻找的答案。我使用了 joda 时间戳,但这不起作用。感谢您的帮助

标签: java apache-spark apache-spark-sql spark-cassandra-connector


【解决方案1】:
java.sql.Timestamp

这用于使用包含时间部分的格式

java.sql.Date

这只是日期

【讨论】:

    猜你喜欢
    • 2020-10-24
    • 2016-08-14
    • 2016-01-05
    • 2020-08-02
    • 2016-09-14
    • 2018-04-03
    • 2015-09-13
    • 2021-02-07
    • 1970-01-01
    相关资源
    最近更新 更多