小编pau*_*ult的帖子

Spark Dataframe - Python - 计算字符串中的子串

我有一个包含字符串类型的列(“assigned_products”)的 Spark 数据框,其中包含如下值:

"POWER BI PRO+Power BI (free)+AUDIO CONFERENCING+OFFICE 365 ENTERPRISE E5 WITHOUT AUDIO CONFERENCING"
Run Code Online (Sandbox Code Playgroud)

我想计算"+"字符串中的出现次数并在新列中返回该值。

我尝试了以下操作,但一直返回错误。

from pyspark.sql.functions import col
DF.withColumn('Number_Products_Assigned', col("assigned_products").count("+"))
Run Code Online (Sandbox Code Playgroud)

我正在运行 Apache Spark 2.3.1 的群集上的 Azure Databricks 中运行我的代码。

python string apache-spark pyspark

10
推荐指数
2
解决办法
6946
查看次数

如何在PySpark中爆炸?

假设我有一个DataFrame用户列和另一列用于他们写的单词:

Row(user='Bob', word='hello')
Row(user='Bob', word='world')
Row(user='Mary', word='Have')
Row(user='Mary', word='a')
Row(user='Mary', word='nice')
Row(user='Mary', word='day')
Run Code Online (Sandbox Code Playgroud)

我想将word列聚合成一个向量:

Row(user='Bob', words=['hello','world'])
Row(user='Mary', words=['Have','a','nice','day'])
Run Code Online (Sandbox Code Playgroud)

似乎我不能使用任何Sparks分组函数,因为它们期望后续的聚合步骤.我的用例是我想将这些数据提供给Word2Vec不使用其他Spark聚合.

apache-spark apache-spark-sql pyspark

9
推荐指数
4
解决办法
7110
查看次数

将组计数列添加到PySpark数据帧

由于其优越的Spark处理能力,我来自R和tidyverse到PySpark,我正在努力将某些概念从一个上下文映射到另一个上下文.

特别是,假设我有一个如下所示的数据集

x | y
--+--
a | 5
a | 8
a | 7
b | 1
Run Code Online (Sandbox Code Playgroud)

我想添加一个包含每个x值的行数的列,如下所示:

x | y | n
--+---+---
a | 5 | 3
a | 8 | 3
a | 7 | 3
b | 1 | 1
Run Code Online (Sandbox Code Playgroud)

在dplyr中,我只想说:

import(tidyverse)

df <- read_csv("...")
df %>%
    group_by(x) %>%
    mutate(n = n()) %>%
    ungroup()
Run Code Online (Sandbox Code Playgroud)

就是这样.如果我想按行数总结一下,我可以在PySpark中做一些简单的事情:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col

spark = SparkSession.builder.getOrCreate()

spark.read.csv("...") \
    .groupBy(col("x")) …
Run Code Online (Sandbox Code Playgroud)

dplyr apache-spark pyspark

9
推荐指数
3
解决办法
1万
查看次数

TypeError:列不可迭代-如何遍历ArrayType()?

考虑以下DataFrame:

+------+-----------------------+
|type  |names                  |
+------+-----------------------+
|person|[john, sam, jane]      |
|pet   |[whiskers, rover, fido]|
+------+-----------------------+
Run Code Online (Sandbox Code Playgroud)

可以使用以下代码创建:

import pyspark.sql.functions as f
data = [
    ('person', ['john', 'sam', 'jane']),
    ('pet', ['whiskers', 'rover', 'fido'])
]

df = sqlCtx.createDataFrame(data, ["type", "names"])
df.show(truncate=False)
Run Code Online (Sandbox Code Playgroud)

有没有一种方法可以通过对每个元素应用函数而不使用?来直接修改ArrayType()列?"names"udf

例如,假设我想将该函数foo应用于"names"列。(我将使用其中的例子foostr.upper只用于说明目的,但我的问题是关于可以应用到一个可迭代的元素任何有效的功能。)

foo = lambda x: x.upper()  # defining it as str.upper as an example
df.withColumn('X', [foo(x) for x in f.col("names")]).show()
Run Code Online (Sandbox Code Playgroud)

TypeError:列不可迭代

我可以使用udf

foo_udf = f.udf(lambda row: [foo(x) …
Run Code Online (Sandbox Code Playgroud)

apache-spark pyspark spark-dataframe pyspark-sql

9
推荐指数
1
解决办法
4438
查看次数

外部覆盖后,Spark和Hive表架构不同步

我遇到的问题是Hive表的架构在使用Spark 2.1.0和Hive 2.1.1的Mapr集群上的Spark和Hive之间不同步.

我需要尝试专门为托管表解决此问题,但可以使用非托管/外部表重现该问题.

步骤概述

  1. saveAsTable一个数据帧保存到一个给定的表.
  2. 使用mode("overwrite").parquet("path/to/table")覆盖数据之前保存表.我实际上是通过Spark和Hive外部的进程修改数据,但这会重现同样的问题.
  3. 使用spark.catalog.refreshTable(...)刷新元
  4. 用表查询表spark.table(...).show().原始数据框和覆盖的数据框之间的任何列都将正确显示新数据,但不会显示仅在新表中的任何列.

db_name = "test_39d3ec9"
table_name = "overwrite_existing"
table_location = "<spark.sql.warehouse.dir>/{}.db/{}".format(db_name, table_name)

qualified_table = "{}.{}".format(db_name, table_name)
spark.sql("CREATE DATABASE IF NOT EXISTS {}".format(db_name))
Run Code Online (Sandbox Code Playgroud)

另存为托管表

existing_df = spark.createDataFrame([(1, 2)])
existing_df.write.mode("overwrite").saveAsTable(table_name)
Run Code Online (Sandbox Code Playgroud)

请注意,使用以下内容保存为非托管表将产生相同的问题:

existing_df.write.mode("overwrite") \
    .option("path", table_location) \
    .saveAsTable(qualified_table)
Run Code Online (Sandbox Code Playgroud)

查看表的内容

spark.table(table_name).show()
+---+---+
| _1| _2|
+---+---+
|  1|  2|
+---+---+
Run Code Online (Sandbox Code Playgroud)

直接覆盖镶木地板文件

new_df = spark.createDataFrame([(3, 4, 5, 6)], ["_4", "_3", "_2", "_1"])
new_df.write.mode("overwrite").parquet(table_location)
Run Code Online (Sandbox Code Playgroud)

使用镶木地板阅读器查看内容,内容显示正确

spark.read.parquet(table_location).show()
+---+---+---+---+
| _4| _3| …
Run Code Online (Sandbox Code Playgroud)

hive mapr apache-spark pyspark

9
推荐指数
1
解决办法
2508
查看次数

如何在 pyspark 中使用 rlike 来使用多个正则表达式模式

我必须使用多种模式来过滤大文件。问题是我不确定使用rlike. 举个例子

df = spark.createDataFrame(
    [
        ('www 17 north gate',),
        ('aaa 45 north gate',),
        ('bbb 56 west gate',),
        ('ccc 56 south gate',),
        ('Michigan gate',),
        ('Statue of Liberty',),
        ('57 adam street',),
        ('19 west main street',),
        ('street burger',)
    ],
    [ 'poi']
)

df.show()
+-------------------+
|                poi|
+-------------------+
|  www 17 north gate|
|  aaa 45 north gate|
|   bbb 56 west gate|
|  ccc 56 south gate|
|      Michigan gate|
|  Statue of Liberty|
|     57 adam street|
|19 west …
Run Code Online (Sandbox Code Playgroud)

apache-spark apache-spark-sql pyspark

9
推荐指数
1
解决办法
2万
查看次数

PySpark:减去两个时间戳列并以分钟为单位返回差异(使用 F.datediff 仅返回一整天)

我有以下示例数据框。date_1 和 date_2 列的数据类型为时间戳。

ID  date_1                      date_2                      date_diff
A   2019-01-09T01:25:00.000Z    2019-01-10T14:00:00.000Z    -1
B   2019-01-12T02:18:00.000Z    2019-01-12T17:00:00.000Z    0
Run Code Online (Sandbox Code Playgroud)

我想在几分钟内找到 date_1 和 date_2 之间的差异

当我使用下面的代码时,它以整数值(天)为我提供 date_diff 列:

df = df.withColumn("date_diff", F.datediff(F.col('date_1'), F.col('date_2')))  
Run Code Online (Sandbox Code Playgroud)

但我想要的是 date_diff 考虑时间戳并给我几分钟的时间。

我该怎么做呢?

python timestamp date apache-spark pyspark

9
推荐指数
1
解决办法
1万
查看次数

如何以两行划分pyspark数据帧

我在Databricks工作.

我有一个包含500行的数据帧,我想创建包含100行的两个数据帧,另一个包含剩余的400行.

+--------------------+----------+
|              userid| eventdate|
+--------------------+----------+
|00518b128fc9459d9...|2017-10-09|
|00976c0b7f2c4c2ca...|2017-12-16|
|00a60fb81aa74f35a...|2017-12-04|
|00f9f7234e2c4bf78...|2017-05-09|
|0146fe6ad7a243c3b...|2017-11-21|
|016567f169c145ddb...|2017-10-16|
|01ccd278777946cb8...|2017-07-05|
Run Code Online (Sandbox Code Playgroud)

我试过以下但是收到错误

df1 = df[:99]
df2 = df[100:499]


TypeError: unexpected item type: <type 'slice'>
Run Code Online (Sandbox Code Playgroud)

python pyspark spark-dataframe databricks

8
推荐指数
3
解决办法
1万
查看次数

多行的 UDF 作为响应 pySpark

我想应用和传递和作为参数的splitUtlisation每一行。,结果将返回多行数据,因此我想创建一个新的 DataFrame (Id, Day, Hour, Minute)utilisationDataFarmestartTimeendTimesplitUtlisation

def splitUtlisation(onDateTime, offDateTime):
    yield onDateTime
    rule = rrule.rrule(rrule.HOURLY, byminute = 0, bysecond = 0, dtstart=offDateTime)
    for result in rule.between(onDateTime, offDateTime):
      yield result
    yield offDateTime


utilisationDataFarme = (
sc.parallelize([
    (10001, "2017-02-12 12:01:40" , "2017-02-12 12:56:32"),
    (10001, "2017-02-13 12:06:32" , "2017-02-15 16:06:32"),
    (10001, "2017-02-16 21:45:56" , "2017-02-21 21:45:56"),
    (10001, "2017-02-21 22:32:41" , "2017-02-25 00:52:50"),
    ]).toDF(["id",  "startTime" ,  "endTime"])
    .withColumn("startTime", col("startTime").cast("timestamp"))
    .withColumn("endTime", col("endTime").cast("timestamp"))
Run Code Online (Sandbox Code Playgroud)

在核心 Python 中,我确实喜欢这样

dayList = ['SUN' , 'MON' …
Run Code Online (Sandbox Code Playgroud)

apache-spark pyspark

8
推荐指数
1
解决办法
4824
查看次数

Spark:日期时间的 col 是工作日还是周末?

我正在用 Python 编写 Spark 代码。我有一个col(execution_date)时间戳。我如何将其转换为名为 的列,如果日期是周末,则is_weekend值为工作日?10

python apache-spark pyspark

8
推荐指数
1
解决办法
5557
查看次数