我想知道 PySpark 中与reset_index()pandas 中使用的命令的等效性。当使用默认命令(reset_index)时,如下:
data.reset_index()
Run Code Online (Sandbox Code Playgroud)
我收到错误:
“DataFrame”对象没有属性“reset_index”错误”
我有一个 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
我正在尝试在 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) 我正在尝试使用 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)。
任何人都知道出了什么问题吗? …
我有一个清单如下:
list=["a","b","c"]
Run Code Online (Sandbox Code Playgroud)
我想转换成以下格式:
data=[{"key": "a"},{"key": "b"},{"key": "c"}]
Run Code Online (Sandbox Code Playgroud)
如何实现这一目标?
我正在 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 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-spark ×5
pyspark ×5
python ×3
apache-kafka ×1
databricks ×1
docker ×1
kafka-topic ×1
postgresql ×1
python-3.x ×1