Apache Spark SQL 是否支持类似于 Oracle 的 MERGE SQL 子句的 MERGE 子句?
MERGE into <table> using (
select * from <table1>
when matched then update...
DELETE WHERE...
when not matched then insert...
)
Run Code Online (Sandbox Code Playgroud) 我正在 databricks 云中运行 pyspark 作业。作为这项工作的一部分,我需要将一些 csv 文件写入数据块文件系统(dbfs),并且我还需要使用一些 dbutils 本机命令,例如,
#mount azure blob to dbfs location
dbutils.fs.mount (source="...",mount_point="/mnt/...",extra_configs="{key:value}")
Run Code Online (Sandbox Code Playgroud)
一旦文件被写入挂载目录,我也试图卸载。但是,当我直接在 pyspark 作业中使用 dbutils 时,它失败了
NameError: name 'dbutils' is not defined
Run Code Online (Sandbox Code Playgroud)
我应该导入任何包以在 pyspark 代码中使用 dbutils 吗?提前致谢。
我正在尝试为 Databricks设置GitHub 集成。
我们那里有数百个笔记本,手动将每个笔记本添加到存储库中会很累。
有没有办法自动提交所有笔记本并将其从数据块推送到存储库?
我收到这个错误
Can't pickle <class 'google.protobuf.pyext._message.CMessage'>: it's not found as google.protobuf.pyext._message.CMessage
当我尝试在 PySpark 中创建 UDF 时。显然,它使用 CloudPickle 来序列化命令,但是,我知道 protobuf 消息包含C++实现,这意味着它不能被腌制。
我试图找到一种方法来覆盖CloudPickleSerializer,但是,我找不到方法。
这是我的示例代码:
from MyProject.Proto import MyProtoMessage
from google.protobuf.json_format import MessageToJson
import pyspark.sql.functions as F
def proto_deserialize(body):
msg = MyProtoMessage()
msg.ParseFromString(body)
return MessageToJson(msg)
from_proto = F.udf(lambda s: proto_deserialize(s))
base.withColumn("content", from_proto(F.col("body")))
Run Code Online (Sandbox Code Playgroud)
提前致谢。
如果满足特定条件,我希望我的 Databricks 笔记本出现故障。现在我正在使用dbutils.notebook.exit(),但它不会导致笔记本失败,并且我会收到类似笔记本运行成功的邮件。如何让我的笔记本出现故障?
我已经安装并配置了 Databricks CLI,但是当我尝试使用它时,我收到一条错误,表明它找不到本地颁发者证书:
$ dbfs ls dbfs:/databricks/cluster_init/
Error: SSLError: HTTPSConnectionPool(host='dbc-12345678-1234.cloud.databricks.com', port=443): Max retries exceeded with url: /api/2.0/dbfs/list?path=dbfs%3A%2Fda
tabricks%2Fcluster_init%2F (Caused by SSLError(SSLCertVerificationError(1, '[SSL: CERTIFICATE_VERIFY_FAILED] certificate verify failed: unable to get local issuer
certificate (_ssl.c:1123)')))
Run Code Online (Sandbox Code Playgroud)
上述错误是否表明我需要安装证书,或者以某种方式配置我的环境,以便它知道如何找到正确的证书?
我的环境是带有 WSL 的 Windows 10 (Ubuntu 20.04)(上面的命令来自 WSL/Ubuntu 命令行)。
Databricks CLI 已安装到 Anaconda 环境中,包括以下证书和 SSL 包:
$ conda list | grep cert
ca-certificates 2020.6.20 hecda079_0 conda-forge
certifi 2020.6.20 py38h32f6830_0 conda-forge
$ conda list | grep ssl
openssl 1.1.1g h516909a_1 conda-forge
pyopenssl 19.1.0 py_1 conda-forge
Run Code Online (Sandbox Code Playgroud)
当我尝试使用 REST …
PyCharm IDE。dbutils.widgets.get()我想在模块中使用,而不是将该模块导入到数据块中。我已经尝试过pip install databricks-client pip install databricks-utils和pip install DBUtils
我想知道是否可以使用代码从笔记本运行 Databricks 作业,以及如何执行
我有一个包含多个任务和许多贡献者的作业,并且我们创建了一个作业来执行这一切,现在我们希望从笔记本运行该作业来测试新功能,而无需在作业中创建新任务,也可以运行循环执行多次作业,例如:
for i in [1,2,3]:
run job with parameter i
Run Code Online (Sandbox Code Playgroud)
问候
首先,我可以说我在写这篇文章时正在学习 DataBricks,所以我想要更简单、更粗糙的解决方案以及更复杂的解决方案。
我正在读取这样的 CSV 文件:
df1 = spark.read.format("csv").option("header", True).load(path_to_csv_file)
Run Code Online (Sandbox Code Playgroud)
然后我将其保存为 Delta Live Table,如下所示:
df1.write.format("delta").save("table_path")
Run Code Online (Sandbox Code Playgroud)
CSV 标题中包含空格和&等字符/,我收到错误:
AnalysisException:在架构的列名称中的“,;{}()\n\t=”中发现无效字符。请通过将表属性“delta.columnMapping.mode”设置为“name”来启用列映射。有关更多详细信息,请参阅https://docs.databricks.com/delta/delta-column-mapping.html 或者您可以使用别名对其进行重命名。
我在该问题上看到的文档解释了如何在使用 创建表后将列映射模式设置为“名称” ALTER TABLE,但没有解释如何在创建时设置它,特别是在使用上面的 DataFrame API 时。有没有办法做到这一点?
有没有更好的方法将 CSV 放入新表中?
更新:
阅读此处和此处的文档,并受到罗伯特回答的启发,我首先尝试了以下操作:
spark.conf.set("spark.databricks.delta.defaults.columnMapping.mode", "name")
Run Code Online (Sandbox Code Playgroud)
仍然没有运气,我遇到了同样的错误。有趣的是,对于初学者来说,将标题中包含空格的 CSV 文件写入 Delta Live Table 是多么困难
我正在研究数据块中的 pyspark。我想生成一个相关热图。假设这是我的数据:
myGraph=spark.createDataFrame([(1.3,2.1,3.0),
(2.5,4.6,3.1),
(6.5,7.2,10.0)],
['col1','col2','col3'])
Run Code Online (Sandbox Code Playgroud)
这是我的代码:
import pyspark
from pyspark.sql import SparkSession
import matplotlib.pyplot as plt
import pandas as pd
import numpy as np
from ggplot import *
from pyspark.ml.feature import VectorAssembler
from pyspark.ml.stat import Correlation
from pyspark.mllib.stat import Statistics
myGraph=spark.createDataFrame([(1.3,2.1,3.0),
(2.5,4.6,3.1),
(6.5,7.2,10.0)],
['col1','col2','col3'])
vector_col = "corr_features"
assembler = VectorAssembler(inputCols=['col1','col2','col3'],
outputCol=vector_col)
myGraph_vector = assembler.transform(myGraph).select(vector_col)
matrix = Correlation.corr(myGraph_vector, vector_col)
matrix.collect()[0]["pearson({})".format(vector_col)].values
Run Code Online (Sandbox Code Playgroud)
直到这里,我才能得到相关矩阵。结果如下:
现在我的问题是:
因为我刚刚研究了pyspark和databricks。ggplot 或 matplotlib 都可以解决我的问题。
databricks ×10
pyspark ×5
apache-spark ×3
python ×3
automation ×1
correlation ×1
ggplot2 ×1
git ×1
github ×1
hadoop ×1
heatmap ×1
jobs ×1
pycharm ×1
pyspark-sql ×1
python-3.x ×1
scala ×1
sql ×1
ssl ×1