代码之家  ›  专栏  ›  技术社区  ›  dejanmarich

将时间戳列与字符串列连接

  •  0
  • dejanmarich  · 技术社区  · 7 年前

    我需要把数据从Cassandra推到ElasticSearch。已从加载数据帧 cassandra ,但列名为 timestamp 是在 Long 格式,所以我需要将其更改为 时间戳 为了更“人类可读”,我使用了:

    val cassDF2 = spark.createDataFrame(rawCass).withColumn("timestamp", ($"timestamp").cast(TimestampType))
    

    数据框现在看起来像:

    +--------------------+--------------------+-------------+--------------------+--------------------+
    |             eventID|           timestamp|       userID|           sessionID|            fullJson|
    +--------------------+--------------------+-------------+--------------------+--------------------+
    |event00001.withSa...| 2018-11-15 09:00...|2512988381908|  WITH_EVENTS_IMPORT|{"header": {"appI...|
    |event00002.withSa...| 2018-11-15 09:00...|2512988381908|WITH_EVENTS_SESSI...|{"body": {}, "hea...|
    |event00003.withPa...| 2018-11-15 09:00...|2006052984315|  WITH_EVENTS_IMPORT|{"header": {"appI...|
    +--------------------+--------------------+-------------+--------------------+--------------------+
    

    现在,我需要连接3列( seesionID, userID and timestamp )变成一个新的( docID )并将其推至ES:

      // concatStrings function
      val concatStrings = udf((userID: String, timestamp: String, eventID: String) => {userID + timestamp + eventID})
    
      // create column docID
      val cassDF = cassDF2.withColumn("docID", concatStrings($"userID", $"timestamp", $"eventID"))
    

    获取错误:

    org.apache.spark.sql.analysiseException:“timestamp”不是数字 列。聚合函数只能应用于数值列。

    我知道 时间戳 是在打电话之后 .cast 现在是一个对象,不能像以前那样聚合 长 ,但如何将其值提取为字符串或其他可以聚合的内容。

    我所能做的就是在 时间戳 列为 长 .

    我的最终数据框架应该是 cassDF2 但是有了新专栏 文档号 其中包含 251929883819082018-12-09T12:25:25.904+0100event00001.withSa... 而不是 15147612000002512988381908event00001.withSa... 在里面 文档号

    1 回复  |  直到 7 年前
        1
  •  2
  •   Leo C    7 年前

    不需要自定义项。您可以使用内置方法 concat 将列组合在一起,包括格式化的字符串 timestamp 具有特定日期格式的列,如下所示:

    import spark.implicits._
    import org.apache.spark.sql.functions._
    import java.sql.Timestamp
    
    val df = Seq(
      ("1001", Timestamp.valueOf("2018-11-15 09:00:00"), "Event1"),
      ("1002", Timestamp.valueOf("2018-11-16 10:30:00"), "Event2")
    ).toDF("userID", "timestamp", "eventID")
    
    val dateFormat = "yyyy-MM-dd'T'HH:mm:ss.SSSZ"
    
    df.
      withColumn("docID", concat($"userID", date_format($"timestamp", dateFormat), $"eventID")).
      show(false)
    // +------+-------------------+-------+--------------------------------------+
    // |userID|timestamp          |eventID|docID                                 |
    // +------+-------------------+-------+--------------------------------------+
    // |1001  |2018-11-15 09:00:00|Event1 |10012018-11-15T09:00:00.000-0800Event1|
    // |1002  |2018-11-16 10:30:00|Event2 |10022018-11-16T10:30:00.000-0800Event2|
    // +------+-------------------+-------+--------------------------------------+
    
    推荐文章