小编Ron*_*ain的帖子

如何在Python中使用Pyspark等效的reset_index()

我想知道 PySpark 中与reset_index()pandas 中使用的命令的等效性。当使用默认命令(reset_index)时,如下:

data.reset_index()
Run Code Online (Sandbox Code Playgroud)

我收到错误:

“DataFrame”对象没有属性“reset_index”错误”

python python-3.x apache-spark apache-spark-sql pyspark

7
推荐指数
1
解决办法
5946
查看次数

如何将sql输出转换为Dataframe?

我有一个 Dataframe,从中创建一个临时视图以运行 sql 查询。经过几次 sql 查询后,我想将 sql 查询的输出转换为新的 Dataframe。我希望将数据放回到 Dataframe 中的原因是这样我可以将其保存到 Blob 存储中。

那么,问题是:将 sql 查询输出转换为 Dataframe 的正确方法是什么?

这是我到目前为止的代码:

%scala
//read data from Azure blob
...
var df = spark.read.parquet(some_path)

// create temp view
df.createOrReplaceTempView("data_sample")

%sql
//have some sqlqueries, the one below is just an example
SELECT
   date,
   count(*) as cnt
FROM
   data_sample
GROUP BY
   date

//Now I want to have a dataframe  that has the above sql output. How to do that?
Preferably the code would be in python …
Run Code Online (Sandbox Code Playgroud)

apache-spark apache-spark-sql pyspark databricks azure-databricks

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

Pyspark 基于具有列表或集合的多个条件的其他列创建新列

我正在尝试在 pyspark 数据框中创建一个新列。我有如下数据

+------+
|letter|
+------+
|     A|
|     C|
|     A|
|     Z|
|     E|
+------+
Run Code Online (Sandbox Code Playgroud)

我想根据给定的列添加一个新列

+------+-----+
|letter|group|
+------+-----+
|     A|   c1|
|     B|   c1|
|     F|   c2|
|     G|   c2|
|     I|   c3|
+------+-----+
Run Code Online (Sandbox Code Playgroud)

可以有多个类别,有许多单独的字母值(大约 100 个,也包含多个字母)

我已经用 udf 做到了这一点,并且运行良好

from pyspark.sql.functions import udf
from pyspark.sql.types import *

c1 = ['A','B','C','D']
c2 = ['E','F','G','H']
c3 = ['I','J','K','L']
...

def l2c(value):
    if value in c1: return 'c1'
    elif value in c2: return 'c2'
    elif value in c3: return …
Run Code Online (Sandbox Code Playgroud)

python apache-spark apache-spark-sql pyspark

4
推荐指数
1
解决办法
5738
查看次数

docker容器中的Pyspark postgresql数据库连接

我正在尝试使用 docker 容器内的 pyspark 连接到计算机的 localhost:5432 上的 postgres 数据库。为此,我使用 VS 代码。VS code 自动构建并运行容器。这是我的代码:

password = ...
user = ...
url = 'jdbc:postgresql://127.0.0.1:5432/postgres'

    
    spark = SparkSession.builder.config("spark.jars","/opt/spark/jars/postgresql-42.2.5.jar") \
        .appName("PySpark_Postgres_test").getOrCreate()
        
    
df = connector.read.format("jbdc") \
.option("url", url) \
    .option("dbtable", 'chicago_crime') \
        .option("user", user) \
            .option("password", password) \
                .option("driver", "org.postgresql.Driver") \
                    .load()
Run Code Online (Sandbox Code Playgroud)

我不断收到同样的错误:

“调用 o358.load 时发生错误。\n:java.lang.ClassNotFoundException:\n无法找到数据源:jbdc。...

也许网址不正确?

url = 'jdbc:postgresql://127.0.0.1:5432/postgres'
Run Code Online (Sandbox Code Playgroud)

该数据库位于端口5432上,名称为postgres。数据库位于我的本地主机上,但由于我在 docker 容器中工作,我认为正确的方法是输入笔记本电脑的 IP 地址 localhost 127.0.0.1。如果您输入localhost,它将引用您的 docker 容器的 localhost。或者我应该使用IPv4 地址(无线 LAN .. 或 wsl)。

任何人都知道出了什么问题吗? …

postgresql docker apache-spark pyspark

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

如何将字符串列表转换为字典值列表

我有一个清单如下:

list=["a","b","c"]
Run Code Online (Sandbox Code Playgroud)

我想转换成以下格式:

data=[{"key": "a"},{"key": "b"},{"key": "c"}]
Run Code Online (Sandbox Code Playgroud)

如何实现这一目标?

python

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

pyspark 数据框中每列的最大字符串长度

我正在 databricks 中尝试这个。请让我知道需要导入的 pyspark 库以及在 Azure databricks pyspark 中获取以下输出的代码

示例:- 输入数据框:-

|     column1     |    column2    | column3  |  column4  |

| a               | bbbbb         | cc       | >dddddddd |
| >aaaaaaaaaaaaaa | bb            | c        | dddd      |
| aa              | >bbbbbbbbbbbb | >ccccccc | ddddd     |
| aaaaa           | bbbb          | ccc      | d         |
Run Code Online (Sandbox Code Playgroud)

输出数据帧:-

| column  | maxLength |

| column1 |        14 |
| column2 |        12 |
| column3 |         7 |
| column4 | …
Run Code Online (Sandbox Code Playgroud)

apache-spark apache-spark-sql pyspark azure-databricks

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

Kafka 不向其他分区发送消息

Apache Kafka 安装在 Mac(英特尔)上。单一本地生产者和单一本地消费者。创建了 1 个具有 3 个分区和 1 个复制因子的主题:

bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --topic animal --partitions 3 --replication-factor 1
Run Code Online (Sandbox Code Playgroud)

生产者代码:

bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic animal
Run Code Online (Sandbox Code Playgroud)

制作人留言:

>alligator
>crocodile
>tiger
Run Code Online (Sandbox Code Playgroud)

生成消息时(通过生产者控制台手动),所有消息都会进入同一个分区。它们不应该跨分区分布吗?

我尝试过 3 条记录(如上所述),但它们仅发送到 1 个分区。在 tmp/kafka-logs/topic-0/00** 00.log 中检查 topic- 中的其他日志为空。

我尝试过几十条记录,但没有成功。

我什至在“config/server.properties”中增加了默认分区配置(num.partitions=3),但没有成功。

我也尝试过不同的主题,但没有运气。

apache-kafka kafka-producer-api kafka-topic kafka-partition

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