Transcription
Databricks 正在迅速成为现代企业的默认数据和人工智能平台。其估值约为 1340 亿美元。而这些数字背后是一个简单的故事。Databricks 不再需要拼接用于数据工程、数据仓库、机器学习、商业智能的独立工具,而是为您提供了一个统一的数据智能平台,所有这些功能都集成在一起。对吧?这就是为什么包括财富 500 强企业中的很大一部分在内的许多组织都在标准化使用 Databricks,以便能够为其从分析报告、AI 代理到 LLM 驱动的应用程序的所有内容提供支持。
作为一名工程师,我有什么优势呢?这仅仅意味着,如果您能在 Databricks 之上构建真正的解决方案,您在就业市场上的价值就会立即提升。因此,在这个我们共同构建的项目中,到项目结束时,我希望您能够非常自信地说,我已经使用 Databricks 构建了一个端到端的餐厅分析平台,涵盖流式处理、批量处理、AI、仪表板,一切都是端到端的。
那么,让我们深入了解细节,我们将一起构建什么?我们将模拟一个在阿联酋五个不同地点运营的大型印度连锁餐厅,拥有 500 名客户和数千条评论和订单。
这就是完整的架构将要呈现的样子。您将使用两个真实的数据源。第一个是用于客户、菜单、历史订单和评论的 Azure SQL 数据库。然后是用于实时流式订单的 Azure Event Hub,模拟一个实时 POS 系统,每隔几秒钟发送新订单。
在 Databricks 内部,我们将使用 Unity Catalog 和 Medallion 架构设置一个完整的 Lakehouse,包含 Bronze、Silver 和 Gold 层。我们将使用 Lakeflow Connect 进行从 SQL Server 到 Bronze 和 Silver 层的变更数据捕获摄取。因此,在这种情况下,源 SQL Server 中的插入和更新会自动流入 Databricks,而无需编写自定义代码。对吧?这就是它的美妙之处。您无需编写任何自定义代码。
然后,我们将使用 Databricks 的另一个有趣功能,即 Spark 声明式管道,将实时订单从 Event Hub 流式传输到 Bronze 层。对吧?这将模拟一个实时 POS 系统。最后,我们将把这些实时订单与一次性历史加载结合起来,以便我们的下游分析,我们构建的仪表板,能够看到完整的图景。
是的。从那里,我们将使用星型模型将原始数据转换为一个非常强大的数据模型。我们将构建“事实订单”或订单级别分析,“事实订单项”,我们将把订单中的所有项展开到不同的行中,以便我们可以看到项级别的性能,最后是“事实评论”表。
有趣的部分就在这里。对于评论,我们将使用 Mosaic AI 的 AI 查询和模型来运行情感分析,并且直接在 SQL 中进行,对吧?直接在 SQL。因此,您将能够理解情感,对客户评论中的问题进行分类,然后在此基础上,我们将创建 Gold 层,用于每日摘要、客户 360 档案和餐厅级别的评论指标。
最后,您将构建两个生产级别的仪表板。第一个是连锁店绩效仪表板。所以,如果您是餐厅老板,您会想看到整个连锁店的运营情况,对吧?这正是这个仪表板的作用。最后是一个评论洞察仪表板,用于跟踪情感、趋势、问题和近期反馈。
所有这些都保留在 Databricks 中。您的数据、AI 和 BI 在一个平台上。
到本视频结束时,您将能够清晰地解释完整的架构。因此,Event Hub + SQL 作为源,LakeFlow Connect 和流式管道用于摄取,Unity Catalog 中的 Medallion 层,AI 用于转换,工作流用于编排,仪表板用于消费。这正是真实团队今天在 Databricks 上构建的方式。
此项目旨在帮助您为面试和日常工作构建真实的技能。在开始之前,还有最后一件事。在整个系列中,您会看到我称之为“概念库”的东西。每当我们遇到一个重要的 Databricks 功能时,比如 Spark 声明式管道或 Unity Catalog,我们都会暂停并深入研究该内容。您不需要有任何先验经验。如果您从未听说过该功能或从未与该功能一起工作过,那完全没关系。我们将一起建立这些知识,并将其存储在您的“概念库”中,以便您能够自信地在面试或实际讨论中解释这些主题。是的。所以,请坚持看完本视频,现在让我们开始构建吧。
首先,让我带您了解我们将要用于此项目的 the data set,最重要的是,我们如何创建 synthetic data。所以,让您知道,这是关于一家大型印度连锁餐厅,分布在阿联酋的五个地点,其中两个在阿布扎比,两个在迪拜,一个在沙迦。对吧?所以,这是我们将要创建的不同数据集。第一个是您的客户,他们大约有 500 名注册客户。餐厅分布在阿联酋的不同地点,然后您有每个不同餐厅的菜单项。历史订单,基本上是过去 6 个月下的订单。所以,我们将把这些数据导入到 Databricks 中的订单表中。对吧?这是过去 6 个月的订单。最后是客户的评论,其中包含您的评分和客户提供的评论文本。是的。
我已经将项目所需的所有数据放在了 synthetic data/data 文件夹中,您也可以在 GitHub 上找到。如果您点击 Databricks 项目 synthetic data data,这些是您可以直接使用的 CSV 文件,而无需运行任何 synthetic data generation 脚本。对吧?但是,如果您有兴趣了解 synthetic data set 是如何创建的,这里有所有的代码片段。第一个是 SQL DB.py,它基本上展示了您的餐厅数据是如何生成的,就在这里。然后您有一套令人垂涎的印度美食项目,它们作为您的菜单项的一部分被生成。最后是您的客户,然后我们将所有这些保存到 CSV 文件中,并保存在 data 目录中。对吧?然后是您的历史订单,如果您还记得,我们刚才谈到了,这些是过去 6 个月生成的订单,对吧?我们读取了我们刚刚生成的餐厅客户和菜单项,以便订单看起来更真实,对吧?订单来自我数据库中存在的餐厅或客户,对吧?然后我们生成这些历史订单,它们基本上也保存在 data 文件夹中。最后是您的评论,我们使用了一些模板,例如五星级评分、四星级评分等等。对吧?同样,这些是使用这种格式生成的,您有评论 ID、订单 ID、给出评论的客户、评论文本以及评分。对吧?然后这最终保存到 customer reviews。是的。
现在请注意,并非所有订单都会有评论,这就是为什么我们在那里设置了一个百分比,它决定了将有多少订单有评论,而这些订单将被保存在这个 CSV 文件中。是的。最后,我们创建了一个脚本,它将生成我们刚才讨论过的所有数据集。对吧?所以,现在我们要运行这个脚本。我们所做的是,我只需运行 python synthetic data.py,然后我们运行它,它将生成所有订单。对吧?所以我已经有了这个,但它一定是用新数据覆盖了它。对吧?这就是您如何为这个项目生成 synthetic data。
现在我们将设置我们的实时数据源,到本部分结束时,您将拥有准备好被 Databricks 消费的、正在流式传输到 Event Hub 的订单。对吧?所以,首先我们需要一个 Event Hub 命名空间。您可以将其视为 Event Hub 的容器,类似于 Kafka 集群包含主题。对吧?所以,让我们转到 Event Hub,创建一个 Event Hub 命名空间。我将创建一个新的资源组。创建新资源组的原因是,以便我们在此项目中创建的所有资源都将保留在此资源组内,对吧?我们将此命名为 eh-en namespace dbx project。对吧?对于这个,因为我们假设连锁餐厅位于阿联酋,并且开发也在阿联酋进行,我们将选择阿联酋北部,我将选择标准定价层。我不选择基本定价层的原因是我想要启用称为 Kafka 表面(Kafka surface)的功能。对吧?Kafka 表面到底是什么?Event Hub 中的 Kafka 表面基本上是指它对 Apache Kafka 协议的内置支持。对吧?它有什么作用?它允许用户将 Event Hub 视为 Kafka 主题。非常有趣。它允许将 Event Hub 视为 Kafka 主题,而无需管理您的 Kafka Broker。对吧?这对我的意义是,Event Hub 将为我提供一个与 Apache Kafka 兼容的终结点,我可以使用它。对吧?我可以使用它与结构化流式 Kafka 连接器。对吧?我可以使用该终结点来读取该主题或 Event Hub 中的所有数据。是的。所以,让我们快速点击“审查并创建”。
好的。现在 Event Hub 命名空间已部署,我们将在其中创建一个 Event Hub。我将将其命名为 orders。对吧?因为我们所有的流式订单都将进入这个 Kafka 主题。对吧?所有其他内容都是正确的。让我们点击“审查并创建”。太棒了。现在 Event Hub 已在此处创建。现在有一件非常重要的事情要做,我们需要创建两个策略,对吧?两个策略。让我快速解释一下这两个策略到底是什么。所以,基本上我们这里有 Event Hub,对吧?Event Hub,让我稍微放大一点。我们有 Event Hub。现在将有一个生产者(producer)来生成这些订单,并将数据发送到 Event Hub。对吧?将所有这些订单发送到 Event Hub。现在 Event Hub 会问“你是谁?”对吧?为什么我应该接受你发送给我的订单或任何记录?对吧?然后这个人,这个生产者将有一个连接字符串(connection string),对吧?它将有一个连接字符串,通过它它将连接到 Event Hub。它将显示这个连接字符串,表明“好的,是你给了我向 Event Hub 推送数据的权限。”对吧?同样,同样,Databricks 将是那个 Databricks 将是那个从 Event Hub 读取数据的人。对吧?同样,Event Hub 再次会问“为什么我要允许你从我这里读取数据?”对吧?然后 Databricks 将再次显示一个连接字符串,对吧?连接字符串表明“好的,是你给了我从这个主题或这个 Event Hub 读取数据的权限。”对吧?所以,这就是为什么我们需要创建第一个,第一个是发送策略(send policy),第二个是监听策略(listen policy)。是的。所以,让我们先创建发送策略。生产者发送策略。当我点击它时,您会看到一个连接字符串。对吧?所以,这是生产者在发送数据到 Event Hub 时将使用的连接字符串。让我再创建一个策略,它将是 Databricks 读取策略。我将点击监听。现在当 Databricks 从 Event Hub 读取数据时,它将使用这个连接字符串。对吧?以便能够读取数据而不会出现任何权限问题。
现在我们的 Event Hub 已设置好,让我们快速看一下我们的订单生成器。您会看到这段代码,对吧?我们在这里获取相同的餐厅、客户和菜单项,以便订单看起来真实。这段脚本与历史订单生成脚本相同,对吧?这段脚本,对吧?所以,我们现在不是保存订单,而是使用 Event Hub Producer Client。对吧?您会看到我们这里有 Event Hub 连接字符串,它将从 Databricks 获取。抱歉,不是 Databricks,而是来自这里的生产者发送策略。对吧?所以,您需要做的是设置您的本地环境。我使用,然后我将创建一个 Cond 环境,并安装我们在这里列出的所有要求。抱歉,这是我们在这里列出的要求集。我将运行这个脚本。对吧?所以,现在我们只需要在这里获取主连接字符串,将其放入一个 env 文件中。现在这将是一个最佳实践。当然,您也可以直接将其放在这里,但我不会推荐这样做。然后我们获取 Event Hub 名称。对吧?这是我们的 Event Hub 名称。我们转到概述。我们转到这里,这将是 orders。让我简单地将其替换为 orders。这将是连接字符串。是的。这就完成了。现在我们创建一个批次。对吧?在一个批次中,我将只附加一个订单。我们可以附加更多订单,但我为了简单起见,将附加一个订单。当然,我们不想产生很多 Event Hub 成本。对吧?最后,我们通过生产者将此批次发送到 Event Hub。对吧?
现在,让我们继续运行此脚本。好的,我们看到它已经开始将订单流式传输到 Event Hub,并且它将以这些秒的间隔进行。对吧?这是 3 秒,对吧?每 3 秒我将发送一个新订单。现在让我们快速验证一下。我们转到数据浏览器,让我选择最新的位置。这里您可以看到一些记录。对吧?这些记录的格式完全相同。例如,让我看看这个。对吧?让我们看看这个订单 ID。这是我们生产者发送的同一个订单 ID。对吧?所以,这现在是我们的实时数据源。在真实场景中,这将是您的 POS 系统、您的移动应用程序或 Web 应用程序推送的订单事件。对吧?所以,我们将保持此脚本运行,并在需要时,当我们开始在 Databricks 中消费数据时。但好消息是我们已经设置好了 Event Hub 中的流式数据源。是的。
现在我们正在设置我们的第二个数据源,即 Azure SQL 数据库。它将包含您的客户、餐厅、菜单项、历史订单和评论的数据。是的。所以,让我们继续创建它。我们将使用相同的资源组。我们将此数据库命名为 restaurant ops。我们将创建一个新服务器。这将是 SQL server ops。我将把它放在同一个区域。让我们使用 SQL 身份验证。暂时使用这个。这将是密码。当然,不是世界上最好的密码,但暂时使用它。是的。工作负载环境将是开发。让我们看一下计算和存储类型。对吧?这个是每月计算成本。vCore,抱歉,这不是每月。这是每 vCore 秒的计算成本。每 vCore 秒的计算成本。这是每月存储成本。基本是 6.46 美元。但让我们选择这个通用和无服务器的。如果您对项目中使用的资源很谨慎,并且在完成项目后立即将其关闭,我非常确定它不会花费您超过 2 美元,一两美元。对吧?我们将选择本地冗余备份存储。我们选择公共终结点,我们需要启用它为“是”。是的。添加当前客户端 IP 地址。所以,当 SQL Server 数据库创建后,我们希望使用 DataGrip 或任何其他客户端上传数据。对吧?因此,为了让 DataGrip 或我们的机器能够与此数据库通信,您需要将此 IP 地址,我们的 IP 地址列入白名单。对吧?这就是我们将要做的。我们选择“是”,然后我们还需要允许 Azure 服务和资源访问此服务器。所以,Azure Databricks 是另一个 Azure 资源。它应该能够访问,它应该能够从 SQL Server 读取数据。所以,这就是为什么我们将其选择为“是”,表示其他 Azure 服务应该能够访问此服务器。是的。然后我们点击“下一步”。所有这些都还可以。对吧?然后让我们快速点击“审查并创建”。
现在我们的 restaurant ops 数据库已准备就绪。现在我将使用 DataGrip 将本地 CSV 文件中的 synthetic data 上传到 Azure SQL 数据库。这是我的 DataGrip 界面。我的意思是,当然,您可以自由使用任何其他客户端,但我经常使用 DataGrip,它非常有帮助。所以,我将转到数据源 Azure SQL 数据库,让我们复制所有这些凭据。restaurant ops,让我们复制 SQL Server 名称,这将是用户名和密码。我们选择的用户名是这个,密码将是这个。数据库是 restaurant ops,对吧?我也将其命名为 restaurant ops。让我们测试连接。完美。这可以。
好的。现在我将向您展示脚本。您进入 SQL,然后看到 Azure SQL 数据库设置。对吧?所以,这些是表。我们已经多次提到它们。我们将首先创建所有这些表,然后我们将使用所有这些代码启用 CDC 和更改跟踪。对吧?是的。所以,让我们继续创建所有这些表,我们将简单地在 DBO schema 中创建它们。对吧?所以,这是 restaurant ops 数据库,这是 DBO schema。对吧?所以,让我们快速执行一个 ima name 并将其替换为 dbo。让我们继续运行所有这些。对吧?所以,所有这些都已连接。如果我刷新它,我应该看到所有这些表。完美。
现在让我们继续插入那些 CSV 文件。导入这些 CSV 文件,对吧?这是您的客户。我将选择客户。好的,这个映射是有意义的。我将为所有其他文件执行此操作,菜单项。现在,让我们快速验证一下。所以,这个显示 auto C1C2。为此,我们需要转到这里。我们选择第一行是标题,然后来到这里,然后现在让我们看看映射。Restaurant ID 已正确映射。Name 已正确映射。对吧?好的。
现在所有数据都已导入,让我们快速看一下这些表。这看起来不错。这看起来也不错。所有这些看起来都不错,对吧?对吧?让我们通过快速计算记录数来进行一次健全性检查。Select count star from dbo customers。应该是 500。完美。历史订单。这应该是 8000。接近 8000。是的。好的。那么我相信所有这些都能正常工作,我们的 synthetic data 现在已上传到 SQL Server 数据库。
还有一个重要步骤,我们需要启用更改跟踪和 CDC。所以,我们需要安装一个实用程序脚本。对吧?我们需要安装一个我将向您展示的实用程序脚本,就在这里。对吧?您可以在我已指定的此链接上找到此脚本。对吧?所以,我们需要安装此实用程序脚本并为 SQL Server 和 Lakeflow Connect 启用更改跟踪(chain tracking)和 CDC。对吧?所以,让我们继续。让我们选择这个数据库,然后转到 SQL 脚本。运行 SQL 脚本,然后转到 SQL 实用程序脚本。让我们继续运行它。好的,这已成功完成。现在我们要执行所有这些语句。对吧?所以,让我们复制所有这些,然后我将其粘贴在这里。是的。所以,数据库名称是 restaurant ops,所有这些都将是相同的。这将是 DO customers。CDC 启用,这基本上执行一个存储过程,我们还为 Lakeflow 启用更改跟踪。让我们在这里输入我们的用户名和密码。对吧?现在我们可以运行它。让我们继续运行它。好的。我们看到更改跟踪已启用,CDC 已创建。是的。完美。太棒了。
所以,我们现在有了两个数据源。Event Hub 正在流式传输实时订单,Azure SQL 数据库现已为批量注入做好准备。
欢迎来到概念库。我们将揭示 Unity Catalog 是什么。但是,如果您已经熟悉 Unity Catalog,请随时跳过本节。对吧?所以,首先,让我们尝试理解为什么会存在 Unity Catalog。对吧?它解决了哪些问题?是的。它解决的第一个问题是数据蔓延(data sprawl)。数据蔓延是什么意思?对吧?所以,理解是您的数据,现实是您的数据分布在许多不同的存储桶、数据库和工作区中。对吧?所以,您的数据分布在许多不同的存储桶中,例如 S3、ADLS,以及许多不同的数据库和工作区。这是当今世界许多组织的现实。没有一个单一的可信赖的地图来了解我们拥有什么数据以及它确切地位于何处。对吧?因此,这导致数据在不同位置和系统中的不受控制的增长。是的。您可能会问为什么?原因是,如果我不知道我的系统已经存在这个特定的数据,我将创建一个管道,一个摄取管道,并摄取该数据。对吧?该数据可能已经存在于某处,但我无法找到它,因为我没有一个单一的可信赖视图。对吧?所以,Unity Catalog 的作用是将其全部拉入,对吧?将所有这些拉入一个统一的、受管的视图,其中包含您所有工作区和所有位置的表、视图、卷和模型。对吧?
问题二是不一致的访问控制。对吧?所以,不同的系统,如云存储、您的数据库和工作区,对吧?数据库和工作区基本上都有自己的安全模型。您授予这些的权限方式各不相同。对吧?因此,它们是单独配置的。因此,会出现重叠和空白。对吧?因此,一些用户可以看到比他们应该看到的更多的数据,而另一些用户则可以完全绕过控制。对吧?所以,Unity Catalog 的作用是提供一个一致的、集中的基于角色的访问控制。对吧?所以,一个一致的、集中的基于角色的访问控制。所以,应该有一个地方,您可以在其中管理您数据系统中所有这些资源的权限和访问。对吧?
问题三是发现挑战。对吧?所以,数据科学家和分析师在进行任何实际工作之前,花费不成比例的时间来查找和理解正确的数据集。对吧?所以,Unity Catalog 的发现层,它为您提供搜索、标签、血缘、质量信号,旨在提高您的数据和 AI 访问的可观察性。对吧?它旨在提高您的数据和 AI 资产的可观察性,这基本上包括您的表、卷、您的 AI 资产将是您的模型。对吧?这基本上是为了让技术人员能够花更多时间来生成实际的见解。是的。
问题四是合规性挑战。对吧?所以,为了简单地回答“在过去 90 天内谁使用了数据集 X?”,人们必须手动拼接来自不同系统的日志,并且仍然不确定这是否是正确的答案。是的。所以,Unity Catalog 会自动捕获用户级别的审计日志。对吧?所以,基本上谁访问了哪个系统,哪个数据集,跨越不同的工作区,并将它们全部放置在系统表中。对吧?我们可以查询这些系统表来了解谁访问了哪个数据集。我们还可以了解我们执行的计算操作或其他操作的成本。是的。
所以,我们已经看到 Unity Catalog 非常出色地解决了所有这些问题,而 Unity Catalog 的目标是提供一个统一的治理解决方案。对吧?它的目标是提供一个统一的治理解决方案。统一,以便您可以在一个地方管理所有数据和 AI 资产。是的。
所以,让我们了解 Unity Catalog 对象模型。是的。Unity Catalog 有一个三级命名空间。对吧?三级命名空间的意思是,为了访问它的任何对象,对吧?您需要通过三级命名约定,即表,然后是 schema,然后是 catalog。对吧?现在这个表可以是表、视图、卷或函数。对吧?现在一个可能的问题,您可能会问,catalog、schema 和最重要的是 metastore 到底是什么?对吧?
所以,metastore 是一个顶级容器。对吧?所以,这是将保存 Unity Catalog 元数据的顶级容器。对吧?现在当一个工作区被创建时,当您创建一个 Databricks 工作区时,Databricks 会自动创建一个 metastore,如果它在该区域中尚不存在。对吧?所以,您将在一个区域中创建您的工作区。如果该区域已经有一个 metastore,它将使用它。如果没有 metastore,它将创建一个。对吧?这个 metastore,它将存储数据在哪里?它将存储在,例如,如果您使用的是 Azure,它将存储在存储帐户中;如果使用的是 AWS,它将存储在 S3 中。现在,这存储在控制平面中。基本上,这意味着 Databricks 将其存储在其自己的存储帐户中。您不会存储它,除非您进行一些配置。对吧?但默认情况下,Databricks 将在其控制平面自己的存储帐户中存储 metastore 的数据,即元数据。对吧?
现在在其下方,我们将拥有 catalog、schema 和表、视图、卷和函数。非常重要的是要注意,一个区域,一个区域只能有一个 metastore。对吧?这是 Databricks 默认制定的规则。对吧?一个区域只能有一个 metastore。对吧?
现在让我们了解 catalog 和 schema 到底是什么。Catalog 是数据组织的主要单元。对吧?它还代表了您如何逻辑地隔离数据。所以,例如,它可以代表组织单位。所以,假设您有一个供应链部门,对吧?我们有一个供应链或一个营销部门,对吧?这些是同一组织的部门。现在这些可能有不同的 catalog,这就是这些部门隔离他们数据的方式,或者数据也可以在不同的软件开发生命周期中隔离。对吧?我的意思是,您可以有一个开发环境,您可以有一个暂存环境,您将有一个生产环境。对吧?所以,同样,您将有一个 dev catalog,您将所有数据集、表、卷、函数和视图存储在 dev 环境中。所以,它们可以是,例如,营销开发,然后您将有营销暂存和营销生产。对吧?同样,您将有,例如,供应链。所以,基本上我们遵循的命名约定或结构是 BU,即您的业务单元,以及软件开发生命周期。对吧?或者我只是在这里将其命名为环境。对吧?我们可以将其命名为 catalog。对吧?
这里的 schema 可以简单地被认为是数据库。作为数据库。然后这些数据库将包含您的表、视图、卷和函数。对吧?所以,通常一个 schema 代表一个单一用例、一个项目或一个团队沙盒。对吧?所以,如果您有一个营销团队,它可以包含,例如,项目一,项目一,然后所有这些表一、表二。类似地,项目二。对吧?或者,这可能是一个不同的用例,等等。对吧?所以,我希望现在 metastore、catalog 和 schema 都有意义了。
现在是时候设置 Databricks 了。我们将创建一个 Databricks 工作区。对吧?这将在同一个资源组中。我将其命名为 RG D,抱歉,不是 RG,而是 WS workspace dbx project。这将是 Premium。所有这些都应该是可以的。对吧?让我们点击创建按钮。对吧?所以,这将为我们提供一个 Databricks 工作区,并且 Unity Catalog 已经配置好。但是,我们需要创建 schema,即 Medallion 架构的 Bronze、Silver 和 Gold。
好的。Databricks 工作区现已创建,我将启动工作区。让我们快速看一下 catalog,我们看到已经创建了一个 catalog。对吧?它与我们的工作区同名,我们将使用它。我将创建四个 schema。我们也可以通过 SQL 编辑器使用 create schema 语句来创建它们,但目前我将使用 UI。对吧?所以,我将创建一个 00 landing,我将将其放在同一个位置,同一个暂存位置。对吧?让我们继续并增加。
好的。所以,我们所有的 schema 都已创建。我们的 Databricks 环境现已准备就绪。对吧?所以在下一部分,我们将使用 Lakeflow Connect 创建我们的第一个数据管道,以从我们的 SQL Server 摄取数据,这是一个具有变更数据捕获的 SQL 数据库。
事情变得令人兴奋了。让我们快速回顾一下我们到目前为止完成的所有事情。所以,第一,我们看到了流式订单生产者脚本。对吧?这基本上是将所有数据推送到 Event Hub。我们创建了 Event Hub 命名空间和 Event Hub 本身,我们看到所有数据都被推送到 Event Hub。然后我们为所有这些创建了 synthetic data,即您的客户、餐厅、菜单项、历史数据和评论。对吧?我向您提供了两个选项。您可以从 data 文件夹获取数据,或者通过运行脚本重新创建 synthetic data。我们创建了 Azure SQL 数据库,然后我们使用了 DataGrip。对吧?当然,您可以自由使用您自己的客户端来导入所有这些数据到 SQL 数据库。对吧?所以,我们已经完成了直到这里的所有事情。现在,我们将使用 Lakeflow Connect 进行从 SQL Server 的 CDC 摄取。是的。
我们已经将 synthetic data set 摄取到了这个 SQL 数据库中。对吧?所以,我们现在将使用 Lakeflow Connect 将所有这些数据从 SQL 数据库摄取到我们的 Unity Catalog 中。对吧?所以,这些表将是您的历史订单、您的评论、dim customer、dim restaurants 和 dim menu item。是的。
所以,Lakeflow Connect 是 Databricks 管理的摄取服务。基本上,它抽象了连接到源系统、跟踪更改和增量加载数据的复杂性。对吧?因此,对于 SQL Server、Salesforce 或 SharePoint 等常见源,无需自定义代码。
所以,现在让我们继续在 Databricks 上创建它。让我们转到作业和管道。在管道中,我们的源是 SQL Server。我们首先创建一个连接。它将通过用户名和密码进行。让我们将其命名为 con restaurant ops。对吧?让我们看看服务器主机是什么,即这个。用户名是这个。密码是这个密码。好的,连接已创建。现在我们看到一些独特的东西,它有两个部分,即您的摄取管道配置和第二部分是摄取网关配置。现在摄取管道是有意义的。对吧?它们将从源摄取您的数据。但是摄取网关配置是什么?所以,摄取网关实际上是一个专用的连续管道,它充当数据提取引擎,基本上从源数据库本身拉取快照和变更数据捕获。对吧?它所做的是将所有这些存储在 Databricks Unity Catalog 的一个卷中,称为暂存区。是的。然后这最终被馈送到主摄取管道,用于将数据馈送到目标表中进行处理。对吧?所以,将摄取网关视为一个连续管道,它跟踪摄取到暂存区(一个 Databricks Unity Catalog 卷)的更改,然后主管道使用该卷将数据摄取到您的实际目标表中。对吧?
一个非常重要的注意事项是,它运行在经典计算上。是的。让我们继续将其命名为 injection gateway。是的。我们将创建两个管道。第一个管道将数据摄取到您的 Bronze 层,第二个管道将专门将数据摄取到 Silver 层。所以,这将是 pipeline injection bronze。对吧?这将是项目。这将是 Bronze。这个将是 Landing。现在让我们继续创建。点击创建管道并继续。
现在在我们等待网关启动的同时,让我向您展示一些有趣的东西。是的。所以,如果我们转到计算,我们看到它正在尝试启动一个集群。这实际上是一个很大的集群。对吧?它最多可以自动扩展到五个这样的工作节点。对吧?当然,我们希望不惜一切代价避免这种情况,因为我们希望控制成本。所以,我们将创建一个策略。对吧?不幸的是,我们无法通过 UI 将策略附加到 Lakeflow Connect。对吧?所以,让我们在这里创建一个策略,它将是,例如,这将是最小计算策略。它将是自定义的。让我们转到高级。这将是工作节点。工作节点的数量将是固定的,即一个。驱动程序,我们将选择允许使用的驱动程序和工作节点的类型。对吧?所以,这个,所以这个应该是 E4D。和 V4。是的。然后让我们添加另一个,即您的工作节点类型。这将是 F4S。F4S,对吧?顺便说一句,我之所以放入所有这些,是因为这些驱动程序和工作节点类型是 Lakeflow Connect 所接受的。是的。好的,这应该足够了。现在让我们继续创建此策略。让我们添加相同的描述,然后让我们继续创建它。
现在这个策略已创建,我希望将其附加到 Lakeflow Connect 启动的任何计算上。对吧?以便我的最小和最大始终保持在 1 和 1 之间。是的。所以,我需要通过 CLI 来完成这个。为此,我们需要安装 Databricks CLI。对吧?所以我安装了这个版本,您也需要安装它。安装完成后,请继续运行 databricks configure。我们需要获取这个 URL 并将其放在这里。这个个人访问令牌将是我们从这里创建的。生成新令牌,这是为了项目。我将简单地复制它并将其放在这里。我也将其保存在这里。是的。这就完成了。现在我希望能够列出我所有的管道。对吧?所以,让我们快速运行它。所以,这些是已创建的管道。是的,有一个网关被创建了,这是我们刚才创建的管道。是的。所以,我想要管道 ID,即这个网关的 ID,因为我们要将策略附加到网关。网关运行在经典计算上。是的。
现在一个非常重要的事情要注意的是,网关运行在经典计算上,但是,但是您的摄取管道将运行在无服务器计算上。所以,如果我转到上一个,这将运行在无服务器计算上,这就是为什么他们会问您是否要添加无服务器使用策略。对吧?在这里他们不会问您任何这样的问题。对吧?因为摄取网关将运行在经典计算上。所以,现在我们继续,我们有了管道 ID。现在我要向您展示的是,让我 [清嗓子] 从我之前运行过的一些东西中复制粘贴。对吧?所以,我们将更新这个管道。让我们获取网关摄取管道 ID 并将其放在这里。这只是 gateway injection,以及到 catalog 的 catalog。所以,基本上我们正在尝试做的是,我们正在尝试附加我们创建的策略。这是我们在这里创建的策略,即您的最小计算策略。让我们获取策略 ID。将其放在这里并附加到此摄取网关。让我们检查一下 catalog 名称是什么。Catalog 名称是 WS DBX project。连接名称是什么?连接名称是这个。复制它。将其放在这里。我相信这个也会改变。网关存储名称,让我们将其与这个相同。让我们继续运行它。完美。这运行了。现在,让我们回到这里。让我关闭它。让我关闭它。对吧?让我们转到作业和管道。我们看到 injection gateway。让我停止它,然后我们运行一个新版本的 injection gateway 管道,看看我们创建的策略是否已附加到计算上。对吧?现在当我们在作业计算中看不到任何东西时。让我们现在运行它。我们应该看到。完美。所以,现在您看到这个计算附加了最小计算策略。现在让我转到这里,您会看到只有一个工作节点。是的。与我们之前看到的五个工作节点相比。是的。
所以,现在让我们回到我们的管道。让我们回到摄取管道。我们点击编辑管道。我们将选取所有应该摄取到 Bronze 层中的表,它们是历史订单。这是评论。是的。我不想更改名称。让我们继续。这将是 Bronze。保存并继续。好的,验证已完成。我现在不会添加计划。让我们简单地保存并运行管道。
好的,完美。嗯,我们的表已刷新,我们看到有历史订单和评论表。我们的 Bronze 层包含其中两个表,而 Silver 层现在是空的。所以,让我们快速看一下。在我启动无服务器启动器仓库之前,一个重要的注意事项是,让我们快速编辑它以选择 2x 小型模型。对吧?并将此设置为 30 分钟不活动。是的。让我们现在开始。完美。这是我们流入的示例数据。同样,对于评论,让我们也看看数据类型是否看起来不错。所以,这看起来不错。评论文本,完美。这看起来也不错。是的。Decimal 如预期。让我们快速计算一下记录数。所以,我们已经选择了原始数据。Select count star from historical orders。这应该给我 8000。完美。这个我不记得数字是多少了,但希望这也会是一个非空数字。完美。所以,现在我们有了工作的 Bronze 管道,它将数据从我们的 SQL Server 摄取到 Unity Catalog 中的 Bronze 层。
所以,现在我们将创建将数据摄取到 Silver 层的管道。对吧?用于创建 dim 表。是的。所以,首先,请记住关闭这两个管道。因为网关是一个连续的摄取管道。对吧?它将一直运行。所以,请务必注意您已关闭这些。是的。所以,现在让我们创建摄取管道。嗯,我们将使用相同的连接。这将是 pipeline injection silver。这将是 Silver。这将是 gateway injection silver。嗯,我实际上应该将上一个命名为 gateway injection bronze,但创建管道是可以的。这已创建,现在我们正在等待我们的网关启动。对吧?我们将需要重复相同的过程。是的。所以,我们转到这里,我们简单地说列出我所有的管道,我们看到应该有一个 gateway injection。是的。所以,您看到 gateway injection silver。现在我们将使用相同的,我们将使用与之前相同的命令。让我们将其粘贴在这里。是的。嗯,这是 Silver 摄取网关的管道 ID。让我将其粘贴在这里。让我 [清嗓子] 将其粘贴在这里。这里有一个。其他所有内容我相信都会保持不变。策略将保持不变。最小计算策略。让我们继续运行它。完美。这已完成。现在让我取消它。让我关闭它。但是您会在这里看到那两个管道被创建了。对吧?gateway injection silver 和 pipeline injection silver。让我停止它并再次运行它,以便新的策略附加到计算上。
好的,现在已停止,让我再次运行它,我们还将运行我们的 Silver 管道。是的。首先,在我们这样做之前,让我们先进行映射。对吧?餐厅到 dim restaurants,客户到 dim customer,依此类推。所以,我们将等待网关启动,然后进行映射。
好的,摄取网关已初始化。现在这次我只想选择客户、餐厅和菜单项。对吧?对于客户,我希望它是 dim customers。对于菜单项,我希望它是 dim menu item。这些是我的 dim 表。对吧?所以,我们还将查看我们将在 Silver 层中创建的数据模型。所以,现在不用担心。是的。嗯,这将是 dim restaurants。是的。让我们继续创建它。这将是 Silver 层。保存并继续。
所以,它现在已经开始验证管道配置。好的。管道已验证并准备就绪。让我们保存并启动管道。是的。
好的。管道似乎已启动,现在它正在提取所有这些表的数据。完美。嗯,我们有 500 名客户,5 家餐厅和 145 个菜单项。所以,我们已经看到我们的 Silver 层摄取管道现已准备就绪并正在运行。但是我想向您展示另一件有趣的事情。对吧?所以,我想向您展示 CDC 将如何工作。所以,让我们在我们的 SQL Server 数据库中更改一些内容,看看它如何在 Databricks 中反映出来。请记住,在完成摄取后关闭所有管道和网关。对吧?我将停止 Silver 管道。我们将使用 Silver 管道进行 CDC 摄取。但是 Bronze 管道将不再使用。对吧?因为它是一次性数据摄取。除非我们想再次进行完全加载,否则我们不需要运行 Bronze 管道。这个用于历史订单和评论。是的。
好的。所以,现在为了查看 CDC 将如何工作,让我们在我们的表中做一些更改。对吧?所以,让我看看客户看起来怎么样。对吧?所以,我们将更改,比如说,这里的其中一个客户。
嗯,他们所在的城市,对吧?所以,我们假设这个更新客户,将城市设置为,比如说阿布扎比,其中客户 ID 等于这个客户 ID 10,000,对吧?所以,嗯,让我们先看看这里的 SQL 编辑器,然后执行一个 select star from 02 silvertim customers where customer ID equals customer。我们运行这个。我们看到这是 Samuel Taylor,住在迪拜。我们继续运行这个。对。这已经成功执行了,让我快速检查一下。我们看到这个已经改变了。对。现在让我们再做一个更改。那么菜单项的总行数是多少?所以,让我快速添加一个新项目到这个餐厅,项目编号是 999。我们运行这个。好的。这已经插入了。现在如果我计数,应该会得到 146。好的。所以,发生了两件新事情,对吧?嗯,一,我们更新了这位客户的城市。二,我们向一家餐厅添加了一个新项目。所以,让我们现在运行我们的管道,运行我们的银色管道,看看 lakeflow connect 是否能够识别 CDC 更改。是的。所以,用于银色注入的网关仍然启动并运行。现在我们看到这是它正在使用的计算,并且注入网关保持运行非常重要,对吧?因为我们之前了解到,正是它跟踪源端的 CDC 更改,然后将其写入数据砖块卷,然后由我们的银色注入管道拾取。是的。所以,让我们现在运行这个。完美。所以,这里您看到两个已更新的记录。一个是 dim customer,另一个是 menu item。所以,让我们快速看一下这里,我们之前看到这是迪拜,现在应该是阿布扎比,对吧?完美。所以,这很有效。我们也可以计数。Select count star from dim silver dot dim menu items。好的,它显示您有 146 条记录。让我们也快速看看 item ID 等于 item 999 的地方。对?我相信这就是 item ID。好的。所以,我们现在看到这一行。所以,这就是使用 CDC 和 lakeflow connect 的优势。我们不需要编写任何代码来检测任何更改。是的。SQL Server 会跟踪它们。Lake flow 会获取它们。Delta Lake 会合并它们,而且是完全自动的。是的。所以,现在当我们完成从 SQL Server 的 lakeflow 注入时,让我们确保我们暂时关闭了所有管道,以尽可能降低成本。是的,欢迎来到概念库。我们将要介绍 Spark 声明式管道。但是,如果您熟悉 Spark 声明式管道,请随时跳过本节。所以,让我们从第一件也是最重要的事情开始,为什么使用 Spark 声明式管道?它解决了什么问题,对吧?它到底是什么,它解决了什么问题?为什么我不能继续,只是使用笔记本电脑编写我的 Spark 代码,然后使用工作流来编排一切?对?正如已经提到的,有两种创建管道的方法。基本上,您编写 pi spark 或 scala 或 sql 代码,然后我们决定作业顺序、调度、错误处理、集群大小,然后我们将所有笔记本粘合在一起,在编排器中,即数据砖块工作流,并安排我们的作业。对?这是最基本和最广泛使用的方法。对?这基本上给了我们很多控制和灵活性。另一方面,当您使用 Spark 声明式管道时,您声明您的表以及如何计算它们,通常使用 SQL 或 pi spark。对?现在 Spark 会自动做的是,它会自动找出依赖关系图。所以,假设您正在使用一个表,它依赖于其他五个表。您不需要确保笔记本 1 必须在笔记本 2 之前运行。Spark 声明式管道会自动找出依赖关系图。对?它将决定执行顺序和您的并行度。是的。一个非常有趣的事情是,它还会处理重试,如果您正在进行流式处理。它会自动处理检查点和增量更新。对?并且它使用相同的定义管理批处理和流式处理。是的。所以,而不是我们一步一步地做所有事情,就像您在第一种方法中看到的那样,您基本上告诉他们,这些是我的源表,这是派生表的代码,对吧?它会为您运行整个过程。您不需要担心您的依赖关系图、执行顺序、事情将如何增量运行、您将如何管理将要到来的数据更改,对吧?所以,所有这些都将由 Spark 声明式管道完成,而且感觉就像魔法,因为您实际上避免了大量代码,否则如果您没有使用公共的,您就不得不编写这些代码。现在,很多人会问,我如何决定何时使用 Spark 声明式管道,对吧?所以,每当您有这个问题时,请问自己一些事情,对吧?其中之一是,您是否希望平台管理编排、可靠性、更改捕获、缓慢变化维度和检查点,对吧?或者您需要对管道进行低级别的精细控制,对吧?所以,当您的工作主要是数据管道,表进表出,并且您想要更多自动化时,请使用 SDP,对吧?关键词,重要的关键词是自动化而不是精细控制。所以,我将带您了解五个标准,您在决定是否使用 SDP、Spark 声明式管道时应始终考虑。第一个是 SQL 风格的 ETL ELT 转换。对?所以,如果您的管道遵循 SQL 密集型模式,即使它们是用 PIS spark 编写的,SDP 也可能是一个不错的选择,对吧?因为 SDP 在声明式表到表转换方面表现出色,对吧?我们将看到声明式到底意味着什么,声明式仅仅意味着您指定您想做什么,而不是一步一步地指定您想如何做,对吧?所以,接下来我们将看一些例子。所以,现在不要太担心这个,对吧?所以,第一个标准是 SQL 风格的 SQL 密集型 ETL ELT 转换。如果这是您管道的组成部分之一,那么 HDP 可能是一个不错的选择。第二点是您想要自动依赖管理,对吧?您希望平台根据您编写的代码自动排序表。是的。第三点是您想要内置的流式处理和增量处理。基本上,您不想过多地担心水印、处理流式概念或添加额外的逻辑来读取增量数据。是的。所以,声明式管道在流式处理和增量数据处理的检查点方面处理得非常好,对吧?所以,如果这是您拥有的约束之一,即您不想过多地担心流式处理逻辑或增量处理,那么 Spark 声明式管道可能是一个不错的选择,对吧?第四点是 CDC 链数据捕获和 CD 缓慢变化维度模式,对吧?所以,基本上 Spark 声明式管道有一个叫做 apply changes 的东西,它有一个叫做 apply changes 的 API,它消除了数百行自定义 Spark 代码,例如进行类型一、类型二 SEDD 处理、乱序事件,并消除了理解水印和流式语义的需要。对?所以,如果您的代码中有这样的模式,并且您不想过多地担心处理逻辑,编写 CDC 和 HCD 的逻辑,HDP 是一个非常好的选择,对吧?最后,操作简单性优于灵活性。所以,正如我之前所说,您不会有很多灵活性,因为您不会有太多的精细控制,但您的管道会非常简单。您不必担心 CDC、SCD、流式处理、增量处理、依赖管理等问题,对吧?而 HDP、Spark 声明式管道提供了声明式函数,对吧?它提供了声明式函数,正如我之前提到的,它将手动 Spark 和结构化流式处理的数百甚至数千行代码减少到几行。现在,正如我之前承诺的,让我们确切地理解声明式是什么意思,对吧?我们如何在这里定义声明式,我们一直在说 FDP SDP 意味着声明式管道,这些管道中的声明式到底是什么意思?声明式是什么意思?声明式仅仅意味着我们指定“什么”,对吧?我们指定“什么”,而不是过多地担心“如何”,对吧?所以,基本上,假设我们有一个输入,我们有一个输入,我们需要执行 XY Z 转换,这会导致一个输出,对吧?所以,我们基本上说我们需要输入和输出,我们不需要担心执行将如何发生,执行顺序将是什么,依赖关系图将是什么样子,我的增量数据将如何处理,所有这些事情将如何发生。我不需要担心,对吧?所以,让我们通过一个例子来理解这一点。假设我没有使用 Spark 声明式管道,而我的流程看起来像这样,如果我有一个表经过青铜、然后是白银、然后是黄金。是的。并且这在您的笔记本 1、笔记本 2 和笔记本 3 中,如这里所示。所以,这段代码,它驻留在笔记本中。这段代码驻留在第二个笔记本中,我们读取青铜,然后我们产生一个输出,该输出被写入白银,然后在第三个笔记本中,我们读取白银,并且输出被写入黄金表。现在我们有三个不同的笔记本,我们需要使用,比如说数据砖块工作流来编排这个并管理依赖关系。对?这是方法一。现在,假设我们正在使用 Spark 声明式管道。所以,我们需要做的是,我们需要指定这是什么应该是输入,以及输出应该是什么样子。对?所以,对于白银,对于白银,让我们先花点时间先覆盖白银。我所说的是,我需要青铜订单,无论我返回什么都将馈送到这个表中,即白银订单。对?所以,Spark 声明式管道的工作方式是,函数的名称就是表本身的名称,除非您在这里指定一个名称。对?所以,我们基本上指定了这是我想要的,这是我想要的,我需要输入是青铜订单表。所以,这是我的输入,而这一行的返回值是我的输出。同样,对于黄金,我也需要白银订单,这就是我想要的输出。输出只是一个分组和聚合。现在,它的美妙之处在于,我不需要担心笔记本 1 是否需要在笔记本 2 之前运行,或者它是否需要在笔记本 3 之前运行。所有这些我都不需要担心。Spark 将查看这段代码,它将自己创建依赖关系图,其中它将有青铜、白银、黄金,然后它将从一些 Delta 文件读取,对吧?并且这个依赖关系图将自己创建。对?所有这些都将由 Spark 管理,所以您可以看到,这变得非常容易,我们在这里编写的代码量与我们在这里编写的代码量相比已经大大减少了,而且这是一个非常简单的例子,我给了您,如果我给您一个复杂的例子,比如说 CDC 或 SCD,我们将在接下来的内容中介绍,您就会明白这确实非常强大。第二个理解声明式的例子是您的 CDC 和 std,即缓慢变化维度处理。对?所以,假设我们有一个叫做 dim customers 的表,这将是一个 SEDD 类型二表。对?这意味着当客户的某些属性发生变化时,我们基本上会跟踪历史记录。对?所以,假设我的初始位置在印度,现在我搬到了新加坡,我将有两个行,它们会说我的位置是从某个日期开始在印度,到某个日期结束。然后我将有另一行,它会说现在我的位置是新加坡,有效日期是今天,但有效日期是空的,表示它是活动记录。现在,如果我们不使用 Spark 声明式管道来做这件事,我们会初始化一个 CDC 流,实时数据正在流入,然后这是已经存在的 dim customers 表,然后我们必须在这里编写大量的逻辑来处理客户记录的 upserts、deletes 和 updates,然后我们还需要管理生效日期,即开始日期和结束日期,以及每条客户记录是否是当前的。我们还需要处理乱序事件,然后我们还需要实现自定义合并逻辑。对?一旦所有这些都完成,我们最终将流数据写入给定位置。对?所以,您可以看到,直到我们达到一个目标,我们需要做很多步骤。对?现在,如果我们使用 Spark 声明式管道,我们所要做的就是指定目标。我们的目标是 dim customers,而 CDC 数据正在流入的源是这里的 CDC customers。我们想要 upsert 的键,或者我应该说比较的键,例如我的名字应该有两个记录。所以,我们应该比较的键是 customer id。我们应该按 updated at 排序,并且应该遵循的 std 是 std2。对?所以,您可以看到我们遵循,我们编写了这么多代码,而不是这么多代码。对?而且这非常容易理解,非常容易编写。所以,回到我们之前所说的,我们不担心“如何”。事情将如何执行?我们只担心“什么”,对吧?输入是什么,输出是什么。对?所以,这决定了输出。对?所以,这就是声明式的含义。在我们开始使用 Spark 声明式管道之前,让我们回顾一些我们需要理解的关键概念。我将使用这张图,它取自数据砖块文档。对?这是非常有趣且重要的文档。所以,您看到这里的第一个路径是流式路径,第二个路径是批处理路径。对?这是批处理路径,这是流式路径。现在,我们将要探索的第一行是流式路径,我们有一个流式源,这个流式源基本上可以是 Kafka、Event Hub、Kinesis、PubSub 等。对?所以,这是您的源,然后通过流式源进入的新记录将通过流式处理。对?当它通过流式处理时,您有两个选项,它可以是追加流,也可以是自动 CDC 流,但让我们退一步,对吧?这些流到底是什么,对吧?所以,流是 Spark 声明式管道中的一个基础数据处理概念,它支持流式处理和批处理。对?支持流式处理和批处理,它基本上从数据源读取数据,应用用户定义的处理逻辑,并将结果写入目标。对?所以,非常简单,一个流基本上从源读取数据,然后应用某种逻辑。对?然后写入目标。对?所以,我们有两种流。追加流和自动 CDC 流。顾名思义,追加流非常简单。它基本上是一个流式流,它读取流式源中的新记录,并将它们追加到目标。对?而不更新或删除之前的行,对吧?从流式源获取行,将它们追加到目标而不更新或删除它。对?[清嗓子] 所以,这就是追加流。那么什么是自动 CDC 流?自动 CDC 流是一种特殊的流式流,它理解 CDC 事件,对吧?它理解更改数据捕获事件,并且它会自动维护缓慢变化维度,对吧?所以,如果我们希望我们的表是 std 类型一或 std 类型二,自动 CDC 流就会理解这一点,它会确保我们的输出表遵循相关的 std,对吧?遵循相关的 st。所以,好处是我们不必担心乱序事件或编写复杂的逻辑,复杂的流式逻辑,以确保 CDC 正常工作,或者缓慢变化维度类型一或类型二正常工作。对?它可以进入两种目标。第一个是您的 synth,这是一个通用的流式目标,例如您的 Delta 表、Kafka 主题、Event Hub 或其他输出。对?而流式表是 Unity Catalog 管理的表,它由一个或多个流式流持续更新。对?所以,可以有一个或多个流式流更新这个流式表。是的。最后,我们有批处理路径,其中我们有一个批处理源。现在,这个批处理源可以是存储中的文件、现有表或在某个时间点物化的流式表。对?现在,来自批处理源的这些数据通过物化流,物化视图流。对?所以,物化视图流是另一个有趣的组件,其中您的数据是增量更新的。对?所以,它只重新处理来自源的更改数据中的新数据,而不是重新计算所有内容。是的。所以,我注意到一个非常重要且有趣的事情是,它并不总是增量的。所以,Spark 所做的是,它会尝试让所有这些逻辑通过一个成本模型,对吧?在通过成本模型之后,它会决定是增量运行还是完全计算。对? wherever possible,它会尝试进行增量运行以避免成本并使其性能优化。对?然后,所有这些数据最终都写入这个物化视图。是的。而物化视图又是 Unity Catalog 管理的表,其内容由查询定义。所以,在这个项目中,我们将同时使用流式表。对?我们将同时使用流式表和物化视图。是的。好的。所以,我希望您现在理解 Spark 声明式管道,并且能够更有意义地理解项目。所以,我们已经从 SQL Server 获得了批处理数据。现在,我们将处理实时数据。是的。我们将创建一个 Spark 声明式管道,它将从 Event Hub 消耗订单并将其写入 Unity Catalog 中的青铜表。是的。一旦完成,我们将一次性加载我们从 SQL Server 加载的历史订单。是的。所以,让我们快速创建一个 ETL 管道。我们将将其命名为 pipeline injection event hub。让我们从一个空文件开始。我将把它重命名为 eventhub.py。在开始将所有数据从 Event Hub 摄取到 Event Hub 的代码之前,让我们创建一个探索笔记本,在那里我们首先在那里测试我们的代码,然后我们将把代码移到声明式管道。对?让我们将其命名为 event hub。是的。我还将快速创建一个我们将用于这里实验的小集群。我们不需要 Photon 加速。这将是单节点。并且让我们在不活动 30 分钟后终止它。是的。我们看到这每小时将花费我们 75 DBU。所以,在它启动的同时,让我们看看如何将数据从 Event Hub 摄取到 Databricks。如果您查看此参考代码,我们基本上会指定我们基本上会获取 Event Hub 的所有详细信息,然后形成所有 Kafka 参数,然后我们在这里使用它来从 Event Hub 或非常可互换地从 Kafka 主题读取数据。对?所以,类似地,我们在这里看到的内容,对吧?我创建了一些类似的东西,在顶部,我们将配置我们的 Kafka 连接。是的。让我把它放在这里。是的。所以,Event Hub 基本上使用 Kafka 协议。所以,我们可以使用 Spark 的 Kafka 源。是的。引导服务器在这里是 Event Hub 命名空间,而主题的订阅是 Event Hub 本身,对吧?我们之前提到的订单主题。所以,让我快速创建所有这些参数。eish name,最后是连接字符串。对?现在,这是 Databricks 想要从 Event Hub 读取数据。所以,让我们快速到这里。我们转到 Event Hub。我们的 Event Hub 名称是 orders。所以,这将是 orders。这个是命名空间。我们把它放在这里。而连接字符串将是,我们转到共享访问策略,然后我们转到 Databricks 读取策略。我们从这里复制连接字符串,然后我们把它放在这里。是的,我希望我。好的,我们的集群已准备就绪,我们运行它。好的,我们运行它。现在,让我们快速创建一个原始数据帧。DF raw equals park dot read and dot format 将是 Kafka,选项将是好的。它看起来会是这样。我将简单地显示 df raw。是的。所以,现在让我们运行它。所以,在它初始化的时候,让我快速解释一下这些参数,对吧?所以,一个非常重要的参数是您的起始偏移量,对吧?嗯,这基本上设置为 earliest,以便我们从流的开头读取。每个触发器的最大偏移量基本上限制了我们每个批次要处理的消息数量。对?所以,它基本上有助于控制资源使用。是的。最后,您有 fail on data loss,它设置为 true。所以,fail on data loss 是 Spark 结构化流式处理 Kafka 源选项,它控制如果 Spark 检测到 Kafka,它期望的某些 Kafka 数据不再可用时会发生什么。对?所以,基本上如果 Kafka,如果 Spark 期望来自 Kafka 的某些数据但它不再可用读取时会发生什么。对?所以,我们将其设置为 true,这意味着在这种情况下它会失败。是的。所以,现在它显示没有行返回。是的。让我快速运行这个流。所以,这应该是 python event.py,这应该会使数据流入其中。对?所以,您看到图表正在上升。是的。现在您看到了一些值。是的。然后这是来自 orders 主题。我们现在只创建了一个分区。是的。所以,现在让我们做一些事情。让我们解析这个数据帧。是的。所以,让我暂时暂停一下。我也暂停一下。是的。然后让我们解析这个数据帧。所以,DF parse 将是您的 DF row,让我们添加 key 列。这里有一些 key 和 value,对吧?所以,我们只将其视为字符串。所以,我们将获取 key str,这将是 key 的列,我们将将其转换为字符串。对?这似乎是我们没有从 pispark.sql functions import star。让我们运行它。这现在应该消失了。是的。然后我们将对 value 做同样的事情。您的 value str 将是这里的 value。是的。现在,在我们检索数据集的任何操作之前,对吧?让我们打印显示 df bar,基本上我们现在看到 key 是 null,这没关系,但 values 中有所有的订单数据,对吧?我们有所有的订单数据,是的。所以,现在我们将要做的是,让我们快速添加另一列,即 with column,这是我的数据,因为这是 JSON,我将解析这个 JSON。是的。所以,让我们从 JSON 中进行,这将是 value str,我们需要在这里传递一个模式,对吧?这个订单的模式,对吧?而模式将是这样的。模式将是这样的。我还将从 pispark.sql.types import star 导入。让我们运行它。然后这应该可以正常工作。然后我将指定订单模式。是的。所以,让我们再次运行它。好的。所以,现在我们看到一个列,其中我们有 JSON。所以,现在您看到它被正确格式化为 JSON 对象。是的。而这属于 data 字段。对?所以,我们将要做的是,我们只是选择这个数据列中的所有字段。对?然后我们也添加重命名。所以,您看到这里有一个 timestamp 列。所以,让我们将其重命名为 timestamp 到 order timestamp。对?所以,我们所做的是,我们获取了 key 和 value 关键字 null,并将 value 传递给字符串。在将其解析为字符串后,我们基本上得到了一个 JSON。但它是字符串格式。我们通过这里的 from JSON 将其转换为 JSON。而这属于 data 列。现在,我们只是选择 data JSON 中的所有列。对?我们只是重命名 timestamp 为 order timestamp。让我们快速运行它。好的。完美。所以,现在您看到的是这里生成的数据,对吧?所有这些在这里生成的数据。我们看到 order ID、order timestamp、restaurant ID、项目的数组、总金额、支付方式和订单状态。完美。这看起来不错。所以,现在让我们将其创建一个 Spark 声明式管道。对?所以,我们首先要做的是,我们首先要暂停这个,对吧?然后我们从这里获取所有内容。是的。我们还将导入 from pi spark import pipelines as dp。是的。然后我们将获取所有这些并将其放在这里。然后我们也获取 Kafka 选项。现在,一件非常重要的事情是,让我们去设置。嗯,我们将选择添加配置。这将是我的值,这将是键名空间。这将是我们的 orders eh eh dot name,然后这将是我的端点,嗯,它将是 eh dot connection string。对?我将简单地将其更改为 spark.com.get get,这将是 eh。这将是 eh dot name,这将是 spark plot ongetk,而这个是连接字符串。对?所以,现在我们拥有了构建 Kafka 选项所需的所有数据。让我们保存它。是的。然后我们将创建管道,对吧?所以,这将通过一个函数 orders。是的。然后让我们粘贴我们创建的所有代码。是的。所以,我们首先创建一个原始数据帧。现在,在创建了原始数据帧之后,我们创建了解析数据帧。是的,解析数据帧是这个。它还需要一个模式。模式在这里。这是您的模式。是的。我们只是返回它。是的。我们只是返回 df parse 数据帧。对?现在,为了让 Spark 声明式管道理解这是一个表,我们需要添加装饰器 DP table。对?我们将把 name 设置为 orders。是的。我们将添加一些表属性,这是一个字典,您可以指定多个内容,我们将提供 quality equals browns。对?现在,一个有趣的事情是,如果您跳过这部分,对吧?如果您跳过名称,它将采用函数名称 orders,并创建一个同名表。对?所以,让我们快速进行一次试运行。但在我们进行试运行之前,让我们检查一下集群。所以,它使用的计算是无服务器的,对吧?所以,让我们快速进行一次试运行。好的。所以,试运行已成功。现在我们将运行管道。好的。所以,它显示从一开始就有 17 个输出。对?所以,如果我们将其设置为 earliest,对吧?事件中心中有 17 条记录。让我也会去 SQL 编辑器向您展示。所以,这是 select。所以,我们有这个,然后在青铜层,我们应该有 orders。好的。我是在错误的地方创建的吗?好的。我的错。我在默认模式下创建的。这应该在青铜模式下。所以,我将再次运行它。然后让我们执行 select star from bronze orders。让我们也进行一次 count star。是的。好的。所以,管道已完成,我相信我现在应该在青铜层看到一些东西。好的,看,我们有 orders,让我快速做这个。好的。所以,我们将快速附加并运行它。好的,看,您看到 17 个订单,但我也运行这个。是的。好的。所以,现在我们有了订单,我们的订单数据从 Event Hub 以流式模式摄取到一个青铜层表。是的。所以我将点击终止并终止此计算,只是为了确保没有其他计算正在运行。嗯,这是我们创建的笔记本,但我们暂时终止它,因为我们不需要它。我们已经看到我们的历史订单表中大约有 8,000 个订单,对吧?而我们刚刚创建了一个流式订单表,正如我们在这里看到的,对吧?所以,DP.t 创建了一个流式表,而我们基本上想将所有历史订单加载到我们当前的订单表中,也就是我们创建的流式表,对吧?所以,我们正在进行一次性将历史数据转储到我们的实际订单表中。对?所以,让我们快速进行。为此,我们将执行 insert into insert into 01 bronze dot orders。我相信鉴于两个表的模式相同,select star 应该可以正常工作。嗯,让我们快速在执行此操作之前进行一次 count star。对?然后让我们在完成后进行一次 count star。好的。所以,插入的行数为 8,000,这是您的所有历史订单。让我们看看这里的数据。我们看到它看起来很好。是的,它看起来很好。所以,现在这应该有 8,000 条记录。是的。所以,现在这基本上是将您的所有历史记录插入到接收您的实时订单数据的同一表中。是的。我们已经配置了从 Event Hub 的实时订单流式传输,并且合并了历史订单,而 lakeflow connect 正在同步我们的批处理。所以,接下来我们将使用这些数据在白银层进行不同类型的转换。所以,接下来我们将创建白银层。这是原始数据变得可靠、类型化和丰富的地方。对?而在这里,我们将创建我们的数据模型。所以,让我快速带您了解我们将要创建的数据模型。它将包含三个事实表。对?和三个维度表。是的。所以,让我稍微放大一下。所以,我们将有 fact orders、fact order items 和 fact reviews。对?所有订单是什么,那些订单中的项目是什么,以及客户的评论是什么,对吧?而维度表是您的 dim restaurants、dim menu items 和您的 dim customers。是的。所以,让我也快速谈谈这些表之间的关系。所以,一个订单,您在这里看到的订单可以有许多 fact order,可以有许多 order items,很简单。是的,这是第一点。一个订单可以从这里的某个餐厅下单,但一个餐厅可以有许多订单,对吧?而且理所当然。然后您看到一个订单可以属于一个客户,但一个客户可以下许多订单,对吧?同样,一个订单将有一个评论,对吧?如果我下了一个订单,我将为该订单写一个评论。而一个评论只能属于一个订单。是的。然后我们同样有其他关系,表示比如说 dim restaurants 和 dim menu items 之间的关系。是的。所以,一个餐厅可以有几个菜单项,但一个菜单项将属于一个餐厅。是的。所以,这是我们将在白银层创建的整体数据模型,它将在下游的黄金层用于构建我们所有的聚合视图。是的。所以,让我们在白银层创建 fact orders 表。您可能会问我它与我们在青铜层已经拥有的 orders 表有什么不同。所以,我们将添加新的属性,它们将是这个。订单小时、星期几、是否是周末、项目数量,对吧?所以,这将有助于我们进行分析,基本上在黄金层构建 KPI。所以,让我们创建一个 ETL 管道,并将其命名为 pipeline transformation transformation silver。让我们从一个空文件开始。是的。所以我将把它重命名为 fact orders py,我们还将添加一个探索文件夹。我们也将这个笔记本命名为 fact orders。对?所以,在这里我们将进行所有实验。是的。所以,让我先复制并记下我们将要创建的表的模式。对?这是模式,现在我们将读取我们的青铜表,order equals spark read。所以,这将是 spark.t bronze dot orders。是的。我们需要创建 order timestamp,这是时间戳类型。让我们快速检查一下青铜层的数据类型,我们已经摄取了这个表。所以,数据类型是字符串,对吧?数据类型是字符串,但其他都很好。总金额是双精度。所以,现在我们将简单地进行一个 with column,这将是 order timestamp,让我们进行一个 to timestamp。好的,让我快速从 pispark.sql.functions import star 导入,这将是 order timestamp 列的 to timestamp。是的。下一个是 order date,这也很简单,with column order date,这应该是 order timestamp 的 to date 列。是的。下一个是 order hour。这也很简单。这应该是 hour。是的。最后是星期几。星期几会有点不同。所以,这是 day of week,我们在这里所做的是,我们首先检查日期格式,对吧?我们检查日期格式,我们以这种形式进行,让我们使用这种日期格式,星期几,我们以 e e 进行,这基本上会给出全名,即星期一、星期二、完整的单词。是的。然后我们检查它是否是周末。是的。所以,为了检查它是否是周末,我们可以,我们所做的是,我们将把它放在一个 f when 中,对吧?让我们删除这个,然后创建。这非常有帮助,对吧?嗯,我不需要写大量的代码。是的。然后我们已经有了 restaurant ID,customer ID 已经有了,order type 也有了,item count。所以,我们需要计算 item count。但我们在这里看到的是,我们有 items,这是一个字符串。是的。但它包含一个 JSON,一个 JSON 数组。是的。所以,让我再次快速向您展示。嗯,如果我们执行 select star,如果我们看到 items,这基本上是一个数组,而数组中包含一个字符串。对?所以,我们将要做的是,我们将创建另一个列,称为 items parsed,这将是来自 JSON,列名是 items,等等 items,我们需要在这里指定我们的模式,item 模式,对吧?我已经创建了 item 模式,它将是这样的。是的。所以,让我把它放在这里。我们需要从park.sql.types import star 导入。让我们运行它。是的。然后我们在这里传递了 item 模式。然后我们最后对 items parsed 进行计数。对?所以,我们进行一个 with column item count,这将是 items parsed 的计数。对?所以,让我们继续,只是显示 DF orders,或者也许我应该将其重命名为 fact orders。是的。让我们继续运行它。好的。好的,我们似乎遇到了一个错误。但在那之前,我认为应该是 size 而不是 count。好的。所以,让我们看看我们创建的列是否已添加到这里。所以,您看到 order timestamp。嗯,让我们看看数据类型,现在是时间戳。而这已经被整齐地放入了我们在这里共享的结构中,对吧?所以,我们添加了星期几、周末、项目数量。好的。这效果很好。所以,让我们选择我们需要的 fact orders 的最终列集。是的。所以,这将是您的 order ID、order timestamp、order date、order hour、您的星期几、周末、restaurant ID、customer ID、order type、item count、total amount,让我们将 total amount 转换为 decimal 以保持一致性。是的,已完成。然后我们添加 payment method 和 order status。所有这些都不需要。是的,让我们在保存之前最后再运行一次。好的,在输入末尾有一些语法错误,基本上是这里缺少括号。订单小时无法解析。好的。所以,这里我创建了 order date 而不是 order hour。现在让我们再次运行它。好的。所以,我们看到 order date、order hour、星期几、周末、item count、total amount 和 order status。所以,这看起来是一个非常整洁干净的 fact order 表。所以,现在让我们将其移至 Spark 声明式管道,对吧?所以我将把所有这些复制到这里。我们还需要一个导入,即 from pi spark import pipelines as your declarative pipeline。是的。然后我们将创建一个函数,我将在其中放置所有这些。是的。所以,让我们创建一个名为 fact orders 的函数。是的。然后我们必须返回这个 packed orders 表。让我们从这里删除显示。为了使它成为一个流式表,我们必须添加 dp dot table。名称将是 fact orders。是的。因为我在这里选择了模式,我还需要更改模式。这将在白银层。所以,这将是 fact orders。您还记得我们添加了表属性。这将是 quality equals silver。是的。所以,让我们对这个表进行一次试运行。所以,问题出在 fact orders。这需要是 df fact orders。好的。这工作得很好。现在让我们运行管道。所以,我们看到我们已经摄取了大约 8k 订单,对吧?让我们去看看这个。好的。让我们也在这里进行一点计数,即您的青铜订单。这个有多少?8017。当我们执行 count star from 0 to silver dot fact orders 时,它也有相同数量的记录。所以,现在我想向您展示一件非常有趣的事情,那就是 Spark 声明式管道中的数据质量。所以,我们可以直接在这里添加数据质量检查,对吧?到我们的 fact orders 管道,我们可以添加的检查类型基本上是,比如说您有一个表,然后我们对这些表设置期望,一系列期望,如果通过了,那就好了,我们保留该记录。如果失败了,您有三个选项。您可以直接删除,您可以删除最终输出中的行,您可以使流失败,您可以使正在运行的管道失败,或者您只是警告并保留记录。是的。所以,我们这样做的方式基本上是放置类似这样的东西,即 dp.expect,然后我们放置期望的名称和 SQL 表达式。对?所以,它看起来会是这样的。我们将对我们的 fact orders 表进行操作。是的。而您对无效记录的选项是,您可以删除它,我们通过简单地放置 DP.exect 或 drop 来删除,或者我们可以使流失败,只需说 DP.exect 或 fail。对?这里有一些例子。Expect or drop。Expect or fail。所以,让我们说我们将使用 expect all or drop。对?所以,有一个叫做 expect all or drop 的东西。是的。所以,这是 expect all or drop。所以,expect all or drop 的意思是我们有一个期望列表,对吧?所以,所有这些期望都应该通过。如果其中任何一个失败,则该行将被删除。是的。所以,让我们写下所有这些期望。所以,我们将说 dp dot expect all or or drop,这将是一个字典,对吧?或者我最好将其格式化为这样的内容。是的。所以,我们将对比如说我们的 order ID 应该是非空的设置一些期望。所以,第一个期望将是 valid order id,这将是 SQL 中的简单写法,order id is not null。是的。同样,您的 order timestamp 也不应该是空的。对?订单,您的 customer ID,您的 restaurant ID,所有这些都不应该是空的。是的。所以,我们将把所有这些放在这里。所有这些都将放在这里。我们还可以包含哪些其他检查?其他一些是,我们有 item count,对吧?所以,从逻辑上讲,item count 应该大于零。所以,让我们也放一个 valid item count,这将是 item count 应该大于零。还有什么?订单状态。所以,订单状态通常应该是有效的分类数据类型,对吧?所以,比如说它应该是 completed、pending、ready 或其他,对吧?所以,我们将导入,我们将基本上把这个检查放在这里。是的。然后我们也将有一些支付方式的类别。是的。所以,对于支付方式,让我们放这个,即 valid payment method,这将是 payment method,它应该在一个列表中,并且它应该只是 cash、card 或 wallet,对吧?最后,我们有总金额,对吧?同样,总金额也应该大于零,而不是零。所以,这将是 valid total amount,而这应该是 total amount 应该大于零。对?所以,我们定义了多少 1 2 3 4 5 6 7 8 个检查我们的表。对?所以,让我们继续运行它。让我们先进行一次试运行。好的,这工作正常。现在让我们运行我们的管道。所以,您看到这里没有输出,对吧?原因是由于我们已经消耗了所有数据。是的。所以,我将在稍后向您展示输出。是的。然而,一个有趣的事情是,存在已注册的期望,并且您看到我们定义了八个期望,对吧?由于我们的流程中没有记录,所以没有采取任何行动。所以,现在让我们继续创建 fact order items。对?这张表提供了项目级别的粒度。而这张表为什么重要?它对于帮助我们回答诸如“我们最畅销的商品是什么?”、“哪个类别产生的收入最多?”之类的问题很重要。对?所以,我将创建另一个
文件,用于我们的事实订单管道。是的,这将是第一个。然后,我们还创建一个笔记本。是的。这个笔记本将被命名为事实订单项。是的。那么,事实订单项将包含哪些属性或列呢?让我把它粘贴在这里。所以,这些是项。所以,这基本上是在订单 ID 项 ID 级别。是的。让我们从 pispark.sql.functions 导入星号,并导入 pispark.sql.type.types 导入星号。是的。所以,现在我们将读取 bronze orders 表。是的。所以,这将是 df fact orders,我们将读取 bronze orders 表。正确。Park.t01 bronze.orders。现在,我们需要创建以下项,即 item ID、restaurant ID、timestamp、order timestamp 可以直接从 orders 表创建。我们将首先创建它,然后使用 order timestamp 列,它看起来会是这样。Order date 看起来也会是这样。呃,order year 和 order month。我们不想要那个,对吧?所以,这是可以的。是的。现在,我们将创建一个叫做 item past 的东西,我们之前也创建过。正确。所以,这将是 item parts。完美。呃,这必须是这里的模式。让我为你复制粘贴模式。我们之前也使用过这个模式。所以,这是我们的模式。是的,这是你的项目模式,首先,让我们在展开之前快速看看它是如何显示的。所以,我们有一个 JSON,一个 JSON 数组,我们将展开它。这意味着一个订单有一个所有项目组成的数组。现在我们想展开它,以便每个项目成为一行。是的。所以,让我们显示 dact orders。但这里有问题。具体是什么问题?好的。这里缺少一个括号。是的。所以,让我们现在运行它。好的。所以,item pass,这是 item,它是一个字符串。item passed,这是一个数组,正如它正确识别的那样,对吧?所以,现在我们希望这些项目中的每一个都成为一行,而所有其他值都将被复制。正确。所以,我们该怎么做?我们只需使用 with column,因为这是一个单独的项目,我将其命名为 item,它将是 all explode。这是对的,是的。现在让我们再看一次。是的。所以,让我们看看这个订单。这个订单。这个订单有这么多项,对吧?它有 chicken tikka masala、masala chai、butter chicken。我不知道谁会喝 chai 配 chicken tikka masala。但是,您可以看到第一个项目是 chicken tikka masala。然后这是 masala chai,对吧?所以,这意味着 explode 已经奏效了。现在我们有了不同的行。我们为订单中的每个项目都有单独的行。是的。所以,现在让我们继续选择 fact order items 的最终列。这将是你的 order ID。下一个是 item ID。选择 item ID 的方式是 item.item id。正确。这将是 column item.item ID。好的。呃,这很好用。然后,您将拥有 restaurant ID、order timestamp、order date 以及与 item 相关的所有其他属性。正确。您将拥有 item name、item category、item quantity、unit price 和 subtotal。正确。现在,一个非常重要的点是,为了避免浮点精度问题,我们将所有这些转换为 decimal。是的,转换为 decimal。这将是 10,2。让我快速复制一下。我们将把它放在这里。是的。所以,现在看起来不错。是的。让我们再运行一次这个管道。好的,这是你的订单。这是 item ID。然后它正确列出了 item name、category、quantity、unit price 和 subtotal。是的。所以,现在让我们把它放入 fact order item。所以,我们将把它粘贴在这里。然后,我们将添加另一行,从 pi spark 导入 pipelines as dp。让我们放入其他代码。我们将把这个函数命名为 fact order items。让我们把所有代码都放在这里。让我们缩进一下。这将是 read stream,正如我们之前讨论过的。DP read stream。让我们返回它。是的。为了将此初始化为表,作为流表,这将是 DP.table,你的名字将是 fact order items,table properties 将是 quality equals silver。是的。所以,让我们做一个干运行。所以,干运行似乎运行良好。让我们继续运行我们的管道。好的。所以,它已经摄取了大约 24K 条记录。是的。所以,这意味着我们的摄取管道运行正常。让我们看看记录。完美。现在,让我们也将数据质量检查作为 fact order items 表的一部分。正确。我们将以与之前相同的方式进行,这将是 expect all or drop。这将是这样的。所以,我们将首先确保诸如 order ID、item ID、restaurant ID 和其他内容之类的东西。订单时间戳不应为空。所以,我们将放置一个有效的检查。Order ID。它不为空。非常好。和 item ID。我们还将放入 restaurant ID 和 valid auto timestamp。正确。这很好。我们不想添加所有其他非空检查。U 一些其他有趣的检查是 quantity。我们在这里看到的数量当然应该大于零。呃,unit price 应该大于零,subtotal 应该大于零。是的。好的,这可以。现在,让我们快速做一个干运行。好的,我们看到定义了七个期望。如果我再次运行管道,我们将看不到任何订单,因为我们已经消耗了 24K 订单,对吧?所以,现在没有新订单进来。但仍然让我们运行管道,看看期望是否已注册。好的,正如预期的那样,没有输出。我们将看到输出,当我们再次运行管道时。是的。但是,我们看到定义了七个期望。现在是最激动人心的部分,我们将构建具有 AI 驱动的情感分析的事实评论。正确。这就是 Mosaic AI 直接集成到我们的管道中的地方。是的。所以,我们将转到 serving,当我们查看这个时,我们看到有几个模型可供我们使用。是的。所以,如果我点击,比如说这个模型,有一个现成的端点,我可以调用它来获取输出。正确。所以,我将简单地使用 databrick 提供的 AI query 函数。AI query 是一个内置函数,它使用给定的提示调用 AI 模型端点。正确。所以,基本上你调用一个端点,提供一个提示,然后它返回输出。我们将提供的提示将要求模型分析评论并返回一个结构化的 JSON,其中包含两个部分。第一部分是你的情感,然后是问题分类。所以,它获取评论文本,我们在 bronze 层中的 reviews 表。是的。然后,它将使用该文本提供给模型。模型将分析我们的文本并返回两件事。第一件事是你的情感。第二件事是问题分类。所以,这是我刚才说的 AI query 函数。基本上,您可以提供模型的端点和您的请求。所以,让我们看一个快速的例子。所以,这个基本上是一个 databrick metal llama 模型。然后,这是我们提供的提示。正确。然后,这将返回一个字符串,即模型的输出。是的。所以,首先,让我准确地向您展示我们将在 factor views 中构建什么。是的。所以,让我们转到我们的 SQL 编辑器,我将放入 fact orders 中的列,属性。是的。所以,基本上将有一个 analysis JSON,它将包含来自模型的响应。所以,我们将获取评论文本,将其发送给模型。模型将告诉我们,嘿,这是评论的情感。它是积极的、中性的还是消极的。除此之外,我们还将对问题进行分类,对吧?所以,订单是否有交付问题,如果有,那么订单评论中确切提到了什么,导致它感觉像是一个交付问题。原因是什么?同样,对于所有其他问题,对吧?因为如果你从业务角度来看,这就是他们想要查看、深入研究和优化的内容,对吧?进行纠正。所以,看看这个 AI query 函数,我们将要做的是,我们将使用这个表,我们摄取的评论表。所以,这将是类似 select star, AI query,我将在这里放置模型,然后我们将有提示。是的。在我粘贴提示之前,让我将其设置为 01 01 bronze.reviews。是的。让我们粘贴提示。提示将是这样的。提示说什么?提示基本上说分析以下评论,并仅返回一个有效的 JSON 对象,其结构如下。结构是什么?结构基本上包含你的情感。它是积极的、中性的还是消极的。然后,它标记是否有交付问题。如果有,那么评论文本中确切说明的原因是什么,对吧?然后,我们最后附加评论文本,这是 bronze review 表中的列。让我将其重命名为 analysis JSON,这就是我们一直在谈论的。正确。所以,现在为了能够选择使用哪个模型,哪个模型,以及现在,让我选择最便宜的一个。这个是每百万个 token 输入 1 DBU。所以,让我们继续使用它,然后回到 SQL 编辑器。是的。所以,这个作为 analytic JSON 应该在这里。所以,让我们现在运行它。好的。现在我们有了输出,让我们,让我们,让我们看看这个。这将很有趣。所以,我们有 analysis JSON,这是情感积极。它有交付问题吗?不。所以,它没有其他问题,对吧?让我们看看消极的。所以,消极的说什么问题是食物,它是食物质量问题,原因是因为令人失望的甜点和甜拉西上的烧焦边缘。我不知道甜拉西上的烧焦边缘是什么意思,但看起来有些不对劲。是的。所以,现在我们有了 analysis JSON,但这是一个字符串,我们需要转换。为此,我们需要提取这里的键。所以,我们这样做的方式如下。我们将使用一些叫做,好的,首先,让我们把它包装在一个 CTE 中。与模型响应作为,让我们把它拿出来,然后把它放在这里。让我们纠正缩进。是的。而我们将能够读取响应的方式是使用 get JSON object。所以,我们将做类似 select 的事情,比如说,我们选择 review id,然后我们做一个 get json object,我们知道这是一个 JSON 字符串对象。这是一个 JSON 字符串对象,它是一个 JSON,这个美元符号基本上指的是 JSON 的根。JSON 输出,我们说我们想要情感,根处的键情感。这将是一列。然后,这个,issue delivery 将是另一列。让我们先选择这个,然后从 model response 中运行它。让我们运行它,看看输出是什么。是的。好的。所以,我们可以看到 sentiment 和 issue delivery 已经正确出来了。现在,我们将对所有其他字段执行此操作。是的。所以,这将看起来像 issue delivery,get JSON object analysis JSON,这将是 dollar issue delivery reason。这将是 text。所以,这就是为什么我们不会做任何改变。它将是 issue delivery region。同样,我们将对其他字段执行此操作。是的。所以,现在我们有了所有用于分类问题的列。是的。让我们添加其他需要的列,即 review ID,然后是 order ID、customer ID、restaurant ID、你的评分和 review text。我们还应该包含 analysis JSON 对象,以防我们以后需要引用它。最后是你的 review timestamp,对吧?这个评论是什么时候给出的?让我们看看最终输出。所以,现在我们看到了输出。让我快速滚动一下。所以,这很有趣。所以,这是积极的情感。如果我们看看一些问题,比如食物令人失望,甜拉西烧焦了边缘。再次,我不知道那是什么意思,因为,是的,嗯,所以,食物质量有问题,情感是消极的,这是有道理的。然后,这是交付问题原因,包装有一些小问题。所以,这意味着交付问题是真的,但情感是积极的。所以,让我们看看情感。它说固体体验,等等,一切都准备得很好,味道很好,但包装上有一些小问题。是的。所以,这很有趣。这有助于我们了解情感是积极的、中性的还是消极的。但它也突出了问题,对吧?订单的问题。好的。所以,现在我们将使用我们在这里创建的 SQL 查询创建一个流表。是的。所以,让我们继续,让我们回到我们的管道。让我现在添加一个文件,这次是 SQL 文件,这将是 fact reviews。是的。所以我将复制我们在这里写的 SQL 查询。是的。所以,它被复制到这里,我们只是创建一个流表。为此的命令是 create or refresh streaming table。这将是 02 silver.fact reviews。是的。作为,然后所有这些。是的。所以,让我们快速做一个干运行。所以,干运行似乎运行良好。我们还没有定义任何数据质量检查。所以,让我们继续。让我们运行管道。好的。所以,有一个问题,无法创建流表 append flow。好的。所以我相信这应该是一个流表,对吧?这个。现在,让我们再次运行我们的管道。所以,这似乎正在运行。太好了。它已经完成,并且有 78 条评论已进入 fact reviews 表。所以,让我们去这里。让我们创建一个新的查询,让我做一个 select star from silver.fact reviews。是的,让我们运行它。瞧。所以,您有 review ID,它属于哪个订单,哪个客户给出了评论,它属于哪个餐厅,评分,评论文本,analysis JSON,然后是分类,你的情感,以及问题的分类。是的。完美。这看起来不错。让我们向我们的管道添加一些数据质量检查。上次我们使用 Python 添加了它。让我们快速看看如何使用 SQL 来做。所以,我们在这里指定约束名称,它与我们之前指定的 valid order ID 相同,然后我们放置期望。但是,如何指定在收到无效记录时该做什么呢?有两种选择,即删除行或失败更新。正确。所以,让我们看一个例子,这个约束是 valid current page,然后我们把它放在这里,在违规时有一个 drop row。是的。所以,让我们做一些类似的事情,我们将在这里创建它。所以,这将是 constraint valid sentiment,我们将写一个期望 sentiment 应该是积极的、中性的或消极的。正确。在违规时,我们将做什么?在违规时,我们将删除行。是的,这有道理。现在,让我们添加另一个检查,用于评分。所以,这将是 constraint valid rating,这将是 expect rating greater than zero,并在违规时 drop。是的。让我们做一个干运行。所以,干运行已经完成。现在,让我们运行管道。由于我们没有推送任何记录,您将再次看不到任何记录。但期望的注册应该发生。所以,我们看到已经注册了两个期望。是的。所以,silver 层现在已完成。我们现在拥有干净的、类型化的和丰富的数据模型。正确。我们已经构建了一个数据模型。我们有事实表和维度表。事实订单用于订单级别分析,事实订单项用于产品分析,事实评论用于 AI 驱动的情感分析。正确。所以,接下来,我们将把所有这些聚合到 gold 层,用于我们的仪表板洞察。所以,让我们快速回顾一下我们到目前为止完成的事情。所以,我们看到我们的订单现在正在流式传输到这个 bronze 表。我们还一次性加载了历史订单到这个 live orders 表。是的。我们已经摄取了 silver 层中的维度表,即 customer、restaurants 和 menu items。我们还使用 Spark 声明性管道创建了 fact orders、fact order items 和 fact reviews。正确。所以,现在我们将开始创建所有 gold 层中的物化视图。所以,gold 是我们的消费层。基本上是为特定分析用例设计的预聚合表。我们将创建三个表。是的。所以,我们将创建一个每日销售摘要,一个客户 360 个人资料。基本上,对于一个客户,它应该能够告诉你一切,对吧?这就是想法。以及餐厅评论聚合,基本上是分析客户评论的情感,他们是否喜欢食物,客户面临哪些问题等等。正确。这些将为连锁餐厅分析提供支持。是的。所以,呃,让我快速带您了解一下我们将要创建的表。所以,让我们看看第一个,这是你的 d sales summary,即你的每日销售摘要。基本上,按每个订单日期,按特定日期,我们将计算所有这些特征或所有这些指标,即对于特定的一天,总订单数是多少?总收入是多少?平均订单价值是多少?唯一客户是多少?订单来自的唯一餐厅是多少?餐饮订单、外卖订单和送餐订单是多少?正确。所以,让我们继续创建一个新的管道。这将是一个 ETL 管道。我们将这个管道命名为 transform gold。我们将从一个空文件开始。所以,我们将把它重命名为 daily sales summary。我们还将创建一个同名的笔记本。是的。让我们删除所有这些。我将使用我之前使用的同一个集群。让我复制我们需要创建的 KPI 或指标。正确。所以,让我们把它放在这里。是的。我们将读取 orders 表。是的。所以,这将是 DF daily aggregation,我们将读取 spark.read,这实际上不是 spark read,这将是 sparktable,这将是 02_sfact orders。正确。集群已准备就绪。让我们看看 fact orders 中的列。是的。所以,我们需要按 order date 分组。我们需要按 order date 分组,对吧?然后,我们需要计算总订单数。所以,要计算总订单数,我们只需执行 count distinct of your order ID,对吧?所以,让我们继续做一个 egg,我们做一个 count distinct。好的,我们需要运行一个导入。所以,从 pispark.SQL SQL dot functions 导入星号,这将是 count distinct of your order id column fact order 将是 order ID,我们将此别名为 total order。然后,我们有 total revenue。Total revenue 应该只是 total amount 的总和。是的。[鼻息声]所以,这应该是 column total amount 的总和,让我们将其四舍五入到 decimal。是的。这将是你的 total revenue。同样,呃,我们有 average order value。我们只需将其转换为 average。这将是 average order value。是的,这将是 average order value。下一个是你的 unique customers。Unique customers 将与这个 count distinct column customer ID 相同,并将此别名为 your unique customers。然后,我们还想找出 unique restaurants。是的,这将是 restaurant ID,这将别名为 unique restaurants。让我们删除这个。下一个是 dining orders、takeaway orders 和 delivery orders。是的。所以,我们所要做的就是 f.sum。或者实际上我们应该计数而不是 f,因为我说 f 的原因是,有时我所做的是导入 pisparksql functions as f,然后做 f.sum 或,是的,类似的东西。是的。所以,这将是 count distinct,我们将放置一个,当你的 column order type 等于 d in 时。当它是 dine in 时,我们只需说 order id。我们只需说 column order id。正确。否则,我们只需说 none。正确。这有道理吗?并将此别名为 9 in orders。正确。现在,我们将为 takeaway 和另一个,即 delivery orders 做同样的事情。是的。所以,就像这样做,我们只需选择这个作为 takeaway。take away,这将是 delivery,这也将是 delivery。是的,完成了。现在,让我们选择,在我们做任何事情之前,对吧?在我们做任何事情之前,让我们看看它是如何显示的。是的。好的。所以,我们有 orders,total order。让我们排序一下。我们有 total orders、total revenue、average order value、unique customers、unique restaurant、dining orders,以及所有这些。正确。完美。这有道理。让我们做一个最终的 select,以确保我们始终选择相关的列。所以,我们将选择 order date,然后是 your total orders、total revenue、average order value、your unique customers,然后是 your unique restaurant banking orders、takeaway orders 和 delivery orders。正确。让我们最后运行一次。好的。所以,这工作得非常好。现在,让我们把它放入声明性管道。是的。好的。所以,现在我们将把我们所有的代码移到 spark 声明性管道。首先,让我们移动导入。我们需要另一个导入,即 pi spark import pipelines as dp。这次我们将这个函数命名为 d sales summary。让我们从这里复制所有代码,然后放在这里。我们将修复缩进。让我们做一个 return。这应该没问题。现在,一个非常重要的注意事项是,这次我们将创建一个物化视图而不是表。是的。所以,基本上我们之前创建的是流表,它们对于 exactly one semantics 和 append only flow 来说很好。是的。所以,它们非常适合你的事件流、CDC 或 IoT 数据。但是,物化视图非常适合维度和聚合,我们在其中进行聚合,对吧?并且在某个地方我们希望查询的结果被预先计算、缓存并自动保持最新,对吧?随着源数据的不断变化。一个非常重要的点我想从文档中展示的是,park 声明性管道为物化视图提供了增量处理引擎。正确。所以,要使用它,你需要用批处理语义编写你的转换逻辑。正确。然后,引擎将只处理新数据和数据源中的更改。正确。所以,这很棒。是的。所以,出于这个原因,我们将只说 DP.read 而不是 read stream。是的。所以,这次,让我们把我们所有的名字都放在这里。名字将是 03 gold d sales summary。我们将对其进行分区。分区列将是 order date,table properties 将是 your quality equals gold。是的。现在,让我们继续运行这个管道,看看它是如何工作的。我也将把它改为 gold。虽然它不会有区别,因为我已经在这里指定了 gold,对吧?让我们快速做一个干运行。所以,干运行非常顺利。现在,让我们继续运行管道。好的。所以,我们看到它已经完成。我们有 182 条记录,输出应该看起来像我们已经看到的,并且确实看起来像我们之前在笔记本中看到的。让我打开 catalog explorer。这是 gold。瞧。您有您的 D sales summary。是的。我们在这里放入的表属性可以在这里作为属性之一。所以,现在我们已经创建了我们的销售摘要。我们的第一个 gold 表,即销售摘要。是的。接下来,我们将创建 restaurant reviews。是的。这将按餐厅聚合你的评论、情感和评分。是的。它还将为我们即将使用的评论洞察仪表板提供支持。是的。它基本上会显示哪些地点有满意的客户,有多少积极、消极的评论随时间变化,它们的趋势等等。是的。所以,让我先向您展示我们究竟要创建哪些指标。是的。所以,让我们先创建一个新的笔记本,这将是 d restaurant reviews。让我们也确保集群正在运行,并且这些都是指标。正确。所以,对于每个餐厅,我们将存储餐厅名称、城市、总评论数、平均评分、你的评分数,获得 5、4、3、2 和 1 评分的次数。是的。然后是情感积极、情感中性和情感消极计数。是的。并把这个表想象成一个每天都会被覆盖的表,或者我应该说已更改的记录将被写入表中。是的。所以,例如,假设一家餐厅有 100 条评论,明天也有 100 条评论。那么该行将不会被覆盖。但是,假设另一家餐厅有 50 条评论。现在有 55 条评论。那么该餐厅的该行将被新数据覆盖。所以,让我们开始编写代码,让我们编写导入,从 pispark.sql.functions 导入星号。让我们运行它,我们将读取 fact reviews 表。是的。所以,这将是你的,所以首先,让我们计算 review stats。统计数据基本上是这些数字,可以很容易地计算出来,对吧?这些都是数字。所以,这将是 review stats,我们将读取 sparkt table,这将是 02 silver.f fact reviews。是的。这次我们将按 restaurant ID 分组,然后进行聚合。是的。我们将聚合并创建所有这些指标。是的。所以,对于 total reviews,这将是 count distinct。Count distinct。实际上,让我打开所有列。这是你的 factory views。这将是 count distinct。review ID,让我将其别名为 total reviews。是的。平均评分将是 average。所以,它已经提供了这个值,让我们也将其四舍五入到两位小数。正确。现在,我们将计算这家餐厅获得了多少次评分为五分。是的。对于这个,我们将使用 f sum f.实际上不是 f.when,只是 when column rating equals 5 then one else zero。这是正确的。这将是你的别名为 rating five count。同样,对于 four count、three count、two count 和 one count。正确。所以我希望这有道理。基本上,我们正在获取你的行。 wherever there is a rating equals 5, we are marking that as one. And wherever there is a rating other than five, we are marking it putting putting the value zero. And then we are summing up the all of the roads, right? And then after summing up, we create the rating five count. Yeah, that makes I hope that makes sense. Now we are going to compute how many times have we got a positive sentiment, neutral and a negative sentiment. Yeah. And this is also going to use the similar approach when column sentiment is going to be positive then we put one otherwise zero and this is going to be my sentiment positive count. And similarly we are going to do it for neutral and negative count. So these are all of the stats that we need for creating our restaurant reviews table. Yeah. So let's go ahead and then do a display on this table. Okay. Okay, the results are here and we see that restaurant ID, total reviews, average rating and there some case why does why didn't the alias work over here? That should have worked. Okay, I should put it over here after the sum. There should be an alias and probably that the reason why it didn't work. Now, let's go ahead and run this and there you go. So, okay, the same mistake, but I hope you get the idea. Yeah. So, now that we've created all of the metrics, the only thing that is remaining now is all of this is done. We need to add the restaurant name and the city. And where is this information present? So, this information is present in your dim table. the dim restaurants table. Yeah, we have the name and the city. So, let's go ahead and do a left join a left join of the restaurant with this review stats. The reason why I'm doing a left join is because your restaurant it is it is a kind of a master table which has a list of all the restaurants, right? So, some will have all of this data, others won't have it and that is completely fine. Yeah, but I want to have a list of all of the restaurant or restaurant reviews. Yeah. So, this is going to be restaurant reviews and we are going to read in the restaurant table. And for reading in the restaurant table, let's first do df restaurant. This is going to be spark.t02 table 02 silver dot dim restaurant and we do a restaurant dot join and this join is going to be with your review stats on the restaurant ID and this is going to be a left join now let's quickly do a display and see what are the item that we get out of it. So, we've gotten everything, right? We've gotten the name, city, country, opening date, phone number, and then all of the other column that were related to the reviews. So, let's quickly select the column that we need. And these are going to be your restaurant ID. This is going to be your restaurant, your name alias as restaurant name and then your city and then all of the other columns coming over here. But there's an important thing that we need to do and that is we need to coales it with zero. The reason why we want to do this is because when we are doing a left join, there are going to be some restaurants which are not going to have reviews and we simply put a zero or a default number for them. Yeah. So average rating this is going to be like this. And let me also cast this number. Okay, let's not do the casting right now. Uh let's do a coal of all the other columns which is your rating five count, four count, three count, two count and one count. And we also do this for the other sentiment column. Yeah. So now this is done and this looks fine. Let's go ahead and run this finally. So we see now what we've wanted. We have the restaurant name, the city, and then your total reviews, average rating, rating count, and all of this. Right? So now we've created the restaurant reviews table. Now let's move all of this code to Spark declarative pipeline. I'm going to create a new file and this is going to be derscore restaurant reviews and we'll start moving in our code from here. Let's move this from pispark pipeline dp. This is going to be named D restaurant reviews. Let's paste all of our code from here. First, we generate the stats table because this is again going to be a materialized view. Let's use DP. And this is going to be removed. Let's place the other part over here which is this one. This is also going to be dp read and we are going to return this data frame the fresh views. Yeah. And this materialize view the name is simply going to be 03 gold and materialize and restaurant view. Yeah, the table properties is going to be quality and gold. This looks good. Let's go ahead and quickly do a dry run. This seems to run fine. We are good to now kick off a run of this pipeline. Okay, so we see five outputs over here and let's go and have a look the restaurant reviews. We have five restaurants and that is why we see five rows over here and this basically come contains all of the details which is your total review, the average rating, all of these counts, your sentiment count. So we've successfully created our restaurant review table. Okay. Now we'll be creating the final table, final goal table which is customer 360. And this is a very important, very meaningful table. Something that almost every organization should have. Yeah. So it basically provides a single view of each customer with all their metrics. So let's say if you are uh if you are an e-commerce company, you would store things like how many orders have a customer made in the last 30 days, 60 days, 90 days. How many complaints have a customer made in the last 30, 60, 90 days, right? How many orders have a customer placed in their whole lifetime? Right? What is their lifetime spend value? Right? So, basically one row to basically describe the entire profile of a customer. Yeah. So, this enables things like personalization, loyalty program, jone analysis. Yeah. So we are we are going to create these columns or these attributes. And let me let me create um another notebook customer 360 and let's first separate this out. Right? So we are going to separate this out in what kind of statistics are they right? So email uh city and all of this these are coming from the customer dimension table. Yeah even join date is going to come from your customer dimension table. So this is going to come from customer table. Loyalty tier total order lifetime spend average order value. last order date, right? Last order date. This is also going to come from your orders. These are basically order statistics. So, let me write this as order stats. These are your review statistics. These are some favorite metrics. And then finally there's another one which is is at is at risk. I think we removed we remove this. We are not going to calculate this. So we are only going to have these metrics to build the profile of a user. Right? So first of all, what are the thing that we need? For this one, we need the customer div table. For this one, we need the fact orders table. Yeah, for this one, we need the fact reviews table. And for this one, we need the fact order items table, right? So then only we'll be able to figure out which f which is the item which is the favorite and the item uh the restaurant that was used to place most number of orders. So let's go ahead and start to build this piece by piece. So first of all, we're going to create order stats and this is simply going to be it is going to read your orders, right? So this is going to be DF orders.t and 02 silver dot fact orders. Right? So this is going to be DF orders and we are going to group by the customer ID and aggregate and create columns such as let's say total order. Let me first frompark.SQL functions import star and we'll be using the same cluster. So this for calculating the total orders I should just do a countd distinct of order id and this should give me total orders for your total amount. Yeah, for total amount the column that I should use is total amount over here. And let's do a sum of total amount. Total amount. And this is going to be your lifetime spend. Average order value is simply going to be this wrapped in actually not wrapped in but instead of sum we use average and this is going to be average order value right let's also round this to two decimal places to keep it neat and then you have your last order date last order date is simply going to be Max of order order date and this is going to be alias last order date. Right? So we have calculated these four values. Right? Now in order to calculate loyalty tier the simple formula that we are going to use is actually we are going to use after this after the aggregation is done we are going to create a column with where your loyalty tier is going to simply follow this kind of a rule when your lifetime spend when your lifetime spend is greater than equal to something let's say 50,000 or Or let's keep this number small because we don't have a lot of orders. Yeah. Uh we have around 8,000 orders. So per customer this 50,000 number might be very big. So let's keep this to be greater than equal to 5,000. Then it will be platinum. Let this be 2,000 then it is going to be gold. And this is going to be 1,000 then it's going to be silver otherwise is going to be bronze and this is going to be your loyalty tier. Yeah. So these are your order stats. Let's go ahead and display this. So we see that per customer we see the total number of orders, lifetime pend, average order value, last order date and the loyalty tier which is gold, silver and platinum. Yeah. So we've created the first part which is your order stats. Now let's create the review stats. For the review stats, we have to read the reviews table and this is going to be spark.t table 02 silver dot fact reviews and what we need to create essentially is per customer the average rating given by them and the total number of reviews. Yeah. So we just do a group by which is review stats equals df reviews dot group by customer id and then we do an aggregation. This aggregation should give me the total reviews. So the total reviews is going to come through count distinct review ID and this is going to be alias as total reviews and the next one is your average ratings given. So this is going to be let's round it off and this is going to be rating rating and we this as average rating given just to make sure that these are the right columns. Factory view does have review ID and rating. So let's display this. Perfect. So now we have the total number of reviews and the average rating given per customer. So next up, let's create another important and interesting metric which is favorite restaurant of the customer. And the way we are defining favorite restaurant is simply the restaurant to which maximum orders were placed by a customer. And we can use this fact order table. We have the customer ID over here and the restaurant ID over here. Right? So let's go ahead and read this table. Actually, we've already read in this table over here, right? So now we do favorite restaurant and this is going to be orders. And let's group by you group by customer ID and restaurant ID. And we do an aggregation by the order count. Yeah. And let's also order this by We order this by customer ID and order count and let's display this. Okay. So we see that for this customer this person has placed two orders from each of these restaurants. So we can simply just pick up any one of it. Right? for this customer. This person has placed maximum number of orders from this restaurant in Dubai which is five orders. Right? So what we are going to do now is simply we are going to use row number. Yeah. We are going to rank we going to partition by the customer ID and we are going to rank by the order count in descending order. So the one which has the highest number of orders goes on top and then we are going to select that restaurant. Yeah. So let's go ahead and create a column called with column row number and this is going to be row number. over. Okay, perfect. Uh and then we are going to order by order count descending and we are also going to filter by we're going to filter by column where Rn equals 1. Right? We need to import this window which is from pispark.sql dot window import window. So this is fine right now and let's go ahead and have a look right now. So now we see that per restaurant you see that this customer like this restaurant in Dubai where he had placed five orders now we see only one row with respect to this customer. Yeah. So we are simply going to remove this order by and let me just select the customer ID and the restaurant ID. Later on we are going to join this with the dim restaurant table and get the actual restaurant name. Yeah, but this should work completely fine. Now let's go ahead and create the final metric which is favorite item. What is the favorite item of each customer? And we have item level details in fact order items table which is your item ID and item name. But we don't have customer ID to be able to link which customer has placed an order for this item. But if you look at fact order, we do have order ID and your customer ID. Yeah. So in [snorts] order to be able to link that these set of items were placed on order by this customer we need to join fact orders and back order item on order id. Yeah. Now this is a classic example of data modeling where if you would have put customer ID in fact order items you wouldn't need to join two large fact tables right and these tables can be very large in several organizations right. So this has been purposely designed this way to show you the importance of doing data modeling the right way. Right? Doing data modeling in a way such that it suits your analytical needs. Yeah. So for now we are going to simply read in fact order items and this is going to be spar.t table. This is going to be 02 silver dot fact order items and this will be your favorite item df order. Let's join it with fact order items and the join is going to be on order ID and it's going to be an inner join which we don't need to specify and we're going to group by customer ID and the item name and let's aggregate and find the quantity. Right? So we have this column called quantity. We basically want to sum the quantity to understand how many number of items were placed for orders. Right? So this is going to be item quantity. And let's do a quick order by. This order by is going to be your customer ID. item name and descending of item quantity. So here we see that this customer who has this ID customer 10,000 probably the Yeah. So, so this customer who had this ID 10,000, the item that he likes the most is this kulfi over here. And if you look at this customer 10,0001, this like this guy likes garlic naan. I I like it too. So, this is where we can now run a window function. And this is going to be Okay, perfect. So this is going to be the same. It's very similar to what we've already done. We are going to do a row number over partition by customer ID. So we are going to partition by customer ID and then order by the item quantity. The item quantity which is the highest goes on top because we are ordering it from the highest to the
最低。而这将得到一个行号为一。所以我们之所以使用行号而不是排名或密集排名,是因为假设有两个项目具有相同的数量,我们只想选择一个项目,对吧?所以,如果第二个项目具有相同的数量,它将自动被编号为二而不是排名中的一。是的。所以我们将在这里按行号过滤,我们只需选择,是的,我们只需选择您的客户 ID 和商品名称,然后让我们将其命名为“最喜欢的商品”。最喜欢的商品。是的。让我们运行一下。好的。似乎有些问题。好的,我们没有设置过滤条件。完美。所以,我们看到客户编号 10,000 的 Kulfi 是最喜欢的商品,客户编号 10,001 的 Galik Nan 是最喜欢的商品。所以,这正如预期那样工作。现在,让我们合并我们创建的所有统计数据,对吧?而基础是我们的客户表,对吧?因为我们想为我们所有的客户生成一个档案。是的。所以,让我们将这个数据框命名为 C3。这将是您的 DF customers,我意识到我还没有读入客户数据。所以这将是 spark 表和 02 silver dot dim customers。这将是 dim customers,我们将继续与 order stats、review stats、favorite restaurant 和 favorite item 进行左连接。对吧?所以,一旦完成,我们将选择我们需要的相关列。是的。所以,我们需要的相关列是,让我们快速复制并粘贴在这里。让我们将其作为 Python 代码。我们也将注释掉。是的。所以,这将是您的客户 ID。客户 ID 将来自客户客户表本身,即 dim customer 表,然后我们需要客户姓名。客户姓名应该在您的 dim customer 表中。所以,让我这样写。别名客户姓名。然后我们有您的电子邮件。电子邮件将保持不变。然后您有您的城市。最后您有您的电话。我们还连接了吗?是的,我们也连接了。抱歉,我们没有电话。我们已经连接了。是的。所以,这是第一部分。现在我们需要终生总订单花费、平均订单价值和最后订单日期。所以,让我们在这里计算,而不是计算,而是放入订单统计信息。所以,这将是,让我们称之为,这将是总订单列。总订单,如果我们客户没有下任何订单,我们将使用零,并且它将再次别名为总订单,我们也会对其他订单做同样的事情。终生花费、平均订单价值、最后订单日期。所以,对于最后订单日期,如果我们是空的,我们将让它保持为空。是的,我们只是保留最后订单日期。最后,还有什么?只是总订单、终生花费、平均订单价值、最后订单日期,最后一个是您的忠诚度等级。是的。所以我相信忠诚度等级我们已经在这里计算过了。所以,我们直接采用这个。忠诚度等级。现在让我们计算,让我们添加评论统计信息。评论统计信息将是这样的,即您给出的平均评分和您的总评论数。这个也将是 coalesced。这将是列,它将是 L0 zero dot alias 平均评分,让我们也将其转换为 decimal dot cost decimal 10 2。我们这样做的原因是,我们想避免任何浮点精度问题。是的。为了保持一致,我们将只使用 decimal,并且还将对所有其他内容使用 decimal,即您的总订单。总订单不是必需的,但终生花费和平均订单价值。是的。然后最后,让我们添加最喜欢的商品列,即您的餐厅 ID 和最喜欢的商品。实际上,如果我们想要,我相信我告诉过您我们会稍后从餐厅 ID 计算名称,但我认为这里是最好的地方。对于最喜欢的餐厅,让我们与您的 dim restaurant 进行连接。所以,DF restaurants 是 dim restaurants。然后我们与 DF restaurants 进行连接,基于 restaurant ID。这将是我们的内连接。让我们看看名称而不是 restaurant ID。好的,有问题。我认为我放错了列名。所以,这里只是名称。对吧?完美。所以,现在看起来不错了。让我们快速重命名。让我们快速将其重命名为餐厅名称。餐厅名称。是的。让我们再次运行它。好的,看起来不错。现在我将在这里放置餐厅名称。是的。而我说的我们没有计算的最后一列是这一列有风险。90 天无订单。而不是这个,我们将计算另一列,即 is VIP。是的。Is VIP,is VIP 基本意味着您的终生花费大于等于 5,000。是的。所以,让我们快速计算最后一列,这将是当列终生花费大于等于 5,000 时,它为真,否则为假,这将是 is vib。好的,现在让我们运行它。我们将显示您的客户 360。就这样。所以我们有了我们的客户、客户姓名、电子邮件、城市、注册日期、您的忠诚度等级、总订单、终生花费、平均订单价值、最后订单日期、给出的平均评分,对吧?然后我们有总评论数、最喜欢的餐厅名称、最喜欢的商品以及这个人是否是 VIP。是的,完美。所以,现在我们已经创建了最终表,即客户 360。现在我们唯一需要做的就是将客户 360 的代码放入 spark 声明式管道。我将再次为客户 360 创建一个文件。实际上有很多代码需要移动。所以,我们将小心一点。让我们复制它。并且将再次从 pi park 导入管道作为 dp。这将命名为客户 360。是的。让我们保持一致。我也会将其重命名为客户 360。是的。所以,现在让我们先放入所有表,对吧?首先我们将放入 order stats,就是这个。我们创建了 order stats,让我们用 db read 替换它。一旦 order stat 被创建,我们就有了 review stats。这再次将是 dpread。Review stats 已合并。现在我们正在放入每个客户的最喜欢的餐厅,对于这个,我们将再次进行 dp readad。让我们将其移到这里。一旦完成,我们将转到最喜欢的商品,也就是这部分。让我们正确缩进,我们也将其替换为 DP read。让我们删除它。最后,我们进行连接,这将生成我们的主 360 客户 360 表。是的。我们将返回这个客户 360 表。是的。好的。所以,我们将再次用 DP read 替换它。让我们快速检查是否有任何错误。好的,这看起来不错。而且,这也将是一个物化视图。我们将将其命名为 03 gold dot customer 360。表属性将与我们之前所做的类似,即 quality is gold。对吧?现在让我们进行一次试运行。试运行似乎没问题。所以,让我们继续运行管道。我们应该在 gold 层看到一个新表。好的。所以,我们看到客户 360 的输出有 500 行,让我们现在运行它。所以,03 gold,这将是您的 D customers 360。让我们运行它。让我们也点击这里看看发生了什么。完美。所以,现在我们已经完成了所有三个 gold 层表。我们的 gold 层已准备就绪。每日销售摘要、客户 360 档案、餐厅评论指标,所有都已预先计算并优化。在下一部分,我们将把所有这些与 data bricks 工作流协调起来。现在,一个重要的注意事项是,我们将使用每日销售摘要和餐厅评论来构建仪表板,但不是客户 360。但这对于理解客户档案是如何在行业中创建的来说是一个重要的表。现在我们将创建一个工作流,它将包含从注入流式注入到 bronze、silver 和 gold 层创建的整个过程。是的。所以,让我们继续进入作业和管道,然后我们点击作业,这将是工作流每日主。所以,我们将添加另一种任务类型,即您的,所以,让我们,让我们实际上退一步。所以,我们创建了很多管道,这是我们正在处理的管道,对吧?管道转换 gold,管道转换 silver,用于创建我们的数据模型,即事实表和维度表,然后我们有管道 bronze,它基本上将所有数据从您的 SQL Server 注入到 bronze 层。管道 silver,它将所有数据从您的 SQL Server 注入到 silver 层作为维度表,即餐厅、菜单项和客户。然后我们在这里进行流式注入。对吧?所以,现在我们想做的是,我们将把它放在一边。我们将按需运行它,只要需要。我们想做的是,我们将把它拿上来,即从 event hub 注入。将其作为第一阶段,然后每当我们有了数据,我们就将运行 silver 转换和 gold 转换。对吧?所以,让我们选择这个。另一种类型将是管道。这个管道是您的 event hub 管道。我们将保持名称相同,管道注入 event hub。我们将创建此任务。现在,一旦此任务创建,我们就需要另一个管道。这个管道是您的 silver 转换,我们将再次保持名称相同,管道转换 silver。它取决于您的 event hub 注入,并且只有在成功后,silver 管道才会运行。我们创建此任务,最后我们创建 gold 任务。是的。所以,这将是 gold,这将是 gold 管道。并且同样,它取决于 silver,正如预期的那样。所以,现在如果我们运行它,我们不应该得到任何输入,抱歉,输出。是的,但让我们继续运行我们的管道,这将是这里的合成订单。让我们开始吧,让我们运行一段时间。让它只推送一些订单到 event hub,然后我们现在运行这个管道。是的,我们点击查看运行。我现在将暂停它,以免推送大量事件到 event hub,对吧?它将获取我们推送的所有事件,并且只有这些事件应该被处理。是的。所以,这个注入已经完成。我们看到 13 个订单。所以,基本上有 13 个订单是从这里推送的。让我们快速看一下。所以,正如预期的那样,有 13 个订单。chil 管道也已成功。所以,我们看到 fact orders,它处理了 13 个订单,我们从 eventhub 获得的订单,这导致了 35 个订单项,并且因为没有评论,所以它保持为空,输出记录保持为空。是的。让我们去 gold 层看看发生了什么。好的。所以,这仍在进行中。所以,我们看到 gold 管道也已计算完成。一个重要的注意事项是,详细摘要的增量化是增量的。这意味着它专门处理了那些记录。它专门处理了所有已更改的订单日期,并且只处理了那些分区。这就是为什么您看到增量化是增量的。现在对于餐厅评论,因为我们没有推送任何评论,所以增量化是无变化的。没有什么改变,对吧?对于这个特定的表。现在,理解一个非常重要的事情是,为什么 Spark 声明式管道对 D customer 360 执行了完全重新计算?而它对销售摘要和餐厅评论却工作得很好。所以,这是一个非常重要的查询,您应该随身携带,基本上我们查看事件日志以了解管道实际运行时幕后发生的事情,即计划信息。是的。所以,我们看到这是用于详细摘要的。这个是用于,是的,这个是用于客户 360 的。所以,让我们复制这个,将其放入 JSON 格式化器。我们在这里看到的是增量计划被成本模型拒绝。因此,维护类型是完全重新计算。选择了完全重新计算,并且有一些成本,我们一直在谈论的成本。我们已经看到 spark 基本上创建了很多计划,而成本模型是最终决定哪个计划将在集群上执行的。是的。所以,值得思考一下,您应该尝试一下。所以,这个客户 360 表在这里,它目前没有按任何列分区,对吧?所以,假设您按客户 ID 分区,然后您流式传输订单,那些订单将包含某些客户 ID,对吧?所以,在分区之后,您需要查看是否只更新了那些客户 ID,因为档案只会为那些客户 ID 更改,对吧?其他客户的所有档案将保持不变。是的。所以,请将此作为家庭作业,并检查如果您按客户 ID 分区会发生什么。然后您流式传输记录,它们将包含某些客户 ID。您的客户 360 表是否会为那些特定记录更新?类似于您在此处看到的 D sales summary 的增量计算模式。是的。所以,现在我们的编排已经完成。它已经就位。流式注入可以持续运行,我们的 silver 层和 gold 层转换在一个主每日工作流中一起运行。是的。所以,接下来我们将构建我们的餐厅连锁店绩效和评论仪表板。所以,让我们快速回顾一下我们到目前为止完成的所有事情。所以,我们的 brown 层已准备就绪。我们的 silver 层已准备就绪。我们构建了数据模型,对吧?我们构建了数据模型,即您的 fact orders、fact order items 和 fact reviews,所有事实表。我们之前已经看到我们已经注入了维度表,对吧?最后,我们创建了 gold 层聚合表,即您的 D sales summary,对吧?基本上在每日级别,在订单日期级别,我们创建了与我们的销售、订单和所有内容相关的摘要统计信息,对吧?然后我们有一个非常独特的表,称为客户 360,它为我们提供了客户的完整档案。例如,我们从该客户赚了多少钱?客户下了多少订单?是的。最后,我们有餐厅评论表,它将为我们提供评论摘要,即客户随着时间的推移收到的评论类型。对吧?而这个表每次都会被覆盖。对吧?这就是它的创建方式。但您也可以创建一个名为 snapshot date 的新列,然后为该特定快照日期追加数据。对吧?所以,您可以按 snapshot date 分区,然后继续追加。这是一种方法。或者您也可以编写类似 merge into 语句的内容。对吧?merge into 语句,并使用 restaurant ID 作为将用于合并的键,并且只有那些有新评论或新评论数据的餐厅才会被更新,对吧?对于我们的情况,我们将做同样的事情,对吧?所以,我们的 spark 声明式管道将找出新的增量数据,然后只有那些餐厅或只有那些角色将被更新,因为有新数据进来。对吧?所以,我们最终创建了所有好的标签。现在我们将创建仪表板。是的。所以,让我们开始吧。所以,现在我们将进入项目的另一个有趣部分,这次我们将构建两个仪表板,并且不需要安装或需要 Tableau 或 PowerBI。我们将把所有东西都保留在 data bricks 生态系统中。所以,我们将使用 databrick AIBI 仪表板,对吧?所以,将有两组仪表板。第一个是您的餐厅连锁店绩效仪表板,对吧?所以,餐厅老板希望了解整个连锁店的表现如何?整个餐厅群体的表现如何?所以,这就是为什么我们有这些指标,即您的总订单、总收入、平均订单价值、高峰时段分析看起来如何?对吧?所以我将向您展示这些仪表板将是什么样子,最终状态将是什么样子,以便您可以获得更多想法。是的,这是第一个,我们可以按特定日期范围过滤,以便我们了解例如总订单是多少,或者收入是多少在这些特定日期范围内。第二个,而且这确实是一个非常重要的一点,是听听客户对食物和产品的感受。对吧?所以,我们可以按餐厅名称过滤,我们将能够查看不同类型的指标。对吧?评论量随时间如何变化?平均评分是多少?对吧?积极、消极和中性评论的数量是多少?是的。我们还将看到这些评论随时间的趋势。最重要的是,我们想了解客户遇到了什么样的问题。是更多地是交付问题?是份量问题?是定价问题?所以,所有这些都将在仪表板上显示。是的。所以,我们将构建的第一个仪表板是您的连锁店绩效仪表板。正如您在这里看到的,有几个指标与订单、收入、活跃客户、每日销售情况、畅销商品、高峰时段分析有关,基本上显示了需求最高的时段。然后您有按订单类型划分的收入。是送货?是堂食?还是外卖?对吧?对吧?哪种订单类型为我带来了最多的收入?然后最后是按食品类别划分的收入。对吧?哪类食品卖得最好。下一个是您的客户评论仪表板。在这里,我们基本上看到每个餐厅。所以,基本上这是您在这里看到的所有餐厅。但我们将可以选择按餐厅过滤,您可以选择一个餐厅,我们将看到总评论数。平均评分是多少?积极、消极和中性评论是多少?总数是多少?对吧?然后我们还看到随时间的分布。这些餐厅获得了一星、二星和五星的评分是多少?对吧?问题是什么?客户遇到的问题是什么?是交付质量、食品质量还是其他?最后是一些最近的积极和消极评论。对吧?所以我希望您和我一样兴奋。现在,让我们开始构建仪表板。所以,让我们开始构建第一个仪表板。是的。我们将创建的第一个指标将是总订单,对吧?而我们将用于此的数据集是您的 gold 表。这将来自 D sales summary,因为我们已经在这里汇总了所有这些信息。对吧?所以,我们来到这里。这将是 D sales summary,这只是一个数字。对吧?所以,让我们取总订单。现在请记住,这是按订单日期进行的,对吧?每个订单日期您将获得总订单,对吧?这就是为什么我们需要对所有订单的总订单求和,以获得这里的总值。对吧?这就是为什么我们求和,这将是总订单。是的。同样,让我们创建所有其他指标,这将是您的总收入。让我们更改此总收入。下一个将是活跃客户。活跃客户的定义基本上是下过订单的客户。对吧?所以,现在我们发现 dales summary 表中没有此信息。我们没有客户 ID。对吧?所以,我们可以在哪里找到此信息?我们可以在 fact order 表中找到此信息。是的。所以,让我们去获取 fact orders 表。是的。现在让我们更改此 fact order,这将只是我的客户 ID,这将是 count distinct。是的。所以,这是活跃客户。现在我们需要平均订单价值。对吧?所以,让我们快速看看我到底在哪里可以使用哪个表?所以,我们不能使用 dail summary,因为我们没有总信息,对吧?我们有总收入,但我不能,只是做一个平均值,因为那将是错误的值,对吧?所以,让我快速检查 fact orders。在 fact orders 中,我有,因为这是在订单级别,我们可以简单地对总金额求和。抱歉,不是求和,而是在这里求平均值。对吧?所以,这将是您的 av,即平均订单价值。最后,我们将拥有唯一客户。对吧?唯一客户。对于这个,我们基本上想了解我们平台上的客户总数是多少。对于这个,让我们使用 dim customer 表。让我们去使用它。让我们在这里放置一个 dim customer。这将是您的客户 ID。count listing customer ID。完美。您的唯一客户再次是 500。所以,基本上这意味着在我们总共 500 名客户中,所有客户至少下过一个订单,对吧?因为这是来自 packed orders,这是来自 dim customer。是的。所以,让我们稍微格式化一下所有这些。让我们也放置一个日期过滤器,对吧?以便我们可以,我们可以随时间查看这些指标。所以,让我将其命名为日期范围。在我们这样做之前,对吧?在我们在这里放置日期范围之前,我们需要使我们在这里使用的指标的底层数据集,它们需要能够使用日期参数过滤数据。对吧?所以,基本上我们将放置 where join date between,让我们在这里添加一个参数。这个参数将是您的 date range dot min 和 date range dot match。为了使其成为参数,我们只需要指定此符号,然后我们放置日期范围,当我们放置最小值时,它将获取我们在这里指定的第一个值,第一个是 from,对吧?所以,让我将其设置为特定日期,开始和结束日期。让我们运行一下。所以,它基本上过滤了此日期范围内的所有值。是的,让我们对所有其他值也这样做。列名将改变。这将是 order date。让我们也将其更改为日期范围。让我们快速放置一个日期范围。是的,让我们现在运行它。同样,对于最后一个也是一样。实际上,dail summary,我们也使用 order date。这只会保持与 back orders 相同。好的,这行得通。现在当我们选择它时,我们基本上想在这里添加所有这些参数。我们去,我们选择一个日期范围选择器,然后我们添加所有参数。所以,我们想用于参数化的字段是哪个字段?它是 date range 字段,对吧?所以,我们也为 fact orders 做。我们也为 dim customers 做。是的。所以,现在您会注意到,每当我在这里更改日期时,所有数字都会得到调整。对吧?所以,您看到数字已经调整了。所以,参数的工作方式就是这样,我们在这里使用日期范围。现在让我们创建每日销售和平均订单价值图表。对于这个,我们再次使用 detail summary。这将是一个折线图。您的 x 轴将是 order date,这将是每日,而不是每月。对吧?让我们将其命名为 order date。y 轴将是您的数字 1 将是总订单。是的。数字 2,好的,这不会是求和,对吧?因为我们已经按订单日期进行了汇总,而这正是我们的 x 轴,对吧?所以,这将是 none。这将是总订单。是的。让我们再添加一个,即您的每日平均订单价值。所以,这再次将是 none。这将是平均订单价值。是的。让我们也快速更改颜色。这个将是这个。是的。所以,总订单是这个总订单,这个是平均订单价值。嗯,让我们,有没有办法我可以,如果我添加标签。好的,这变得非常混乱。所以,让我们不要这样做。是的。所以,这就是我们的每日销售情况。例如,在 10 月 20 日,我们的总订单是 33,平均订单价值是 171。下一个是畅销商品前 10 名,对吧?而我们有商品级别数据的地方是哪里?商品级别数据存在于 back order items 中。所以,让我们快速运行它。这里有商品名称。要了解销售了多少次,我们取 quantity 列,对吧?所以,让我们只做一个 item name 和 quantity 的总和作为 total quantity sold,然后我们按 item name 对 total quantity sold 进行降序排序,然后我们限制为 10。对吧?但现在一个非常重要的事情是,我们想在日期范围内这样做。对吧?所以,这将是 where order date between date range dot min 和 date range dot max。对吧?让我们将其更改为日期范围。我们将在这里选择日期范围。是的。让我们现在运行此查询。好的,这已根据我们选择的日期进行了修改。现在让我们回到,也让我命名一下。让我们将其重命名为畅销商品前 10 名。这将是这里的一个新图表。让我快速添加,让我们快速添加这个。畅销商品前 10 名必须在日期范围内。这就是为什么我们在这里添加此过滤器,这将来自畅销商品前 10 名。它将是一个条形图。y 轴将是商品名称,x 轴将是 total quantity sold。对吧?所以,我们将删除求和,因为我们已经在查询中进行了求和。对吧?让我删除这个。是的。所以,我们已经进行了求和。让我们也快速更改颜色。所以,这是 total quantity sold。这是您的商品。所以,这里我们创建了一个图表来找出最畅销的 10 种商品。对吧?所以,让我们在这里添加标题。畅销商品前 10 名。接下来,我们要构建的是按星期几划分的订单量。对吧?所以,对于这个,我们将使用 fact orders,x 轴将是您的星期几,我们已经在这里有了。订单量或订单量,我们将简单地选择您的 order ID。订单 ID 在哪里?这将是 count distinct。对吧?所以,这将是您的总订单。总订单。这是您的星期几。对吧?让我们在这里也放一个标签,以便我们能够快速了解在特定工作日下了多少订单。标题将是按星期几划分的订单量。对吧?接下来,我们要进行高峰时段分析。对吧?所以,基本上我们想了解每天在哪个时段需求最高和最低,对吧?所以,我们将为此使用热力图。让我们在这里放置热力图,并使用 fact orders 表,因为它有订单小时,并且还有日期和星期几。对吧?所以,我们将使用热力图在我的 x 轴上,我们将拥有星期几。对吧?实际上,不是在 x 轴上。让我们将其放在 y 轴上。所以,在 y 轴上,我们有星期几。在 x 轴上,让我们放置订单小时。对吧?现在,为了,为了决定颜色,这只是按您的 order ID 的 count distinct。对吧?所以,order ID 的 count distinct。让我们使用红色作为颜色。让我们检查一下。是的,看起来好多了。所以,我们将使用这个。所以,我们看到在星期三的订单,在订单小时 13,即午餐时间左右,订单量很高。对吧?同样,在星期二,午餐时间,大约下午 2 点,订单量很高,对吧?并请注意,这是使用合成数据生成的,对吧?所以,让我们稍微美化一下这里的文字,这将是订单小时,对吧?这将是星期几。是的。最后,让我们构建最后两个图表,即按订单类型划分的收入。所以,让我在这里再放一个图表,按订单类型划分的收入。我认为 fact orders 中有订单类型。所以,这将是一个条形图。在 y 轴上,我们将拥有您的订单类型,在 x 轴上,我们将拥有收入,即您的总金额,这将是求和,因为我们还没有预先汇总。所以,这将是您的总收入,让我们快速在这里写下标题。按订单类型划分的收入。是的。让我们也在这里放置标签。这将是总收入。这是您的订单类型。是的。最后一个是我按食品类别划分的收入是多少?如果我快速在这里放置,按食品类别,我们在 fact order items 中有类别级别的详细信息。所以,我们将选择 fact order items。让我们在这里放置。让我们快速复制这个,以便我们可以进行过滤,其中 order range between min 和 max,这将是一个日期范围。让我们快速选择一个特定日期范围。让我们运行它。好的,这没问题。所以,这次我们将将其制作成饼图。这将是您的颜色。您的颜色将由类别决定,在哪里。好的,我正在看 fact order。这应该是 fact order items。颜色应该由类别决定,角度应该由收入决定。所以,让我们在这里对 subtotal 求和。是的,我们对 subtotal 求和,它看起来是这样的。是的。所以,sum subtotal 将被称为 total revenue,category 将只是 category。让我将其大写。让我们在这里放置一个标题,即按食品类别划分的收入。好的。所以,我们已经构建了第一个仪表板,它看起来非常不错。让我检查一下我们是否为所有使用的数据集都设置了参数。所以,我们使用 factor order items。让我们也为其设置一个参数。对吧?让我们继续发布此仪表板。所以,我们已经在这里发布了这个仪表板。它看起来很棒。是的。让我们也尝试更改日期。让我们放在这里。让我们放 1 月 10 日。我们看到数字已经改变了。是的,正如预期的那样。完美。现在让我们继续进行评论仪表板。所以,这将是客户评论仪表板。我们将使用 fact,抱歉,不是 fact,而是我们为此目的创建的 D restaurant reviews 表。我们看到每个餐厅都有特定的统计信息,对吧?即您的总评论数、平均评分、五星、四星、三星、二星、一星评分的数量,然后是您的积极、消极和中性评论的数量。对吧?所以,这次让我们继续添加这个。我们将创建我们的第一个指标,这只是一个数字。值将是您的总评论数。您看到我们已经对总评论数进行了求和。求和基本上是对连锁店中所有餐厅的所有评论进行求和。对吧?所以,这基本上是我的总评论数。是的。这就是为什么我们放置了求和,因为我们想要连锁店的总评论数。对吧?但现在,如果我们想按特定餐厅过滤呢?是的。我们将在这里使用一个字段。对吧?所以,让我们选择餐厅名称。对吧?所以,这里您看到所有,如果我们选择所有,它应该加总所有评论。是的。我还想与您分享为什么我们使用字段而不是参数。是的。因为之前您已经看到我们在所有地方都使用了参数。所以,参数基本上是占位符,它们在运行时修改底层数据集。对吧?所以,您已经看到我们将其注入到查询本身。例如,其中所有日期介于 date range dominax 之间。所以,这基本上提供了对查询逻辑非常强大的控制。对吧?所以,那时您可以使用参数。而字段过滤器则引用查询运行后数据集的已解析结果上的过滤器。对吧?所以,基本上它们用于提供非常快速的交互性,用于已获取的数据。所以,如果我在这里选择这个餐厅,它将变为 25。对吧?所以,无论是否求和,结果都将相同。是的。您看到,如果我放 none,我也得到 25。如果我放求和,我也得到 25,因为正如我们在表中看到的,每个餐厅只有一行。是的。所以,现在,让我们添加一些其他指标。这将是平均评分。对于这个,我将使其成为平均平均评分,因为同样,如果它是一个餐厅,它没有区别。但当您对多个餐厅进行操作时,它基本上是对已经平均的评分进行平均。对吧?所以,它基本上将所有这些加起来,除以五,然后当您选择所有时,它会给出最终的平均评分。是的。然后我们还想了解这个餐厅在哪个城市。所以,让我们选择城市本身。这将是 none。这将是城市。是的。让我们这样写。然后我们将添加三个新列,即好的,让我们将其缩小。我将删除这个。我们将在这里添加三个新图块,它们将告诉我您的积极评论、您的积极评论数、您的中性评论以及最后您的消极评论。是的。所以,对于这个,我们也需要这个值,即积极计数。是的。我们只需要求和。是的,出于与之前相同的原因。我们取这个是中性,我们求和,对于这个,我们选择消极,我们求和。对吧?让我们也对它进行一些颜色编码。这将是红色。这将是黄色。这将是绿色。是的。接下来,我们将构建一个有趣的图表,即监控评论情绪随时间的变化趋势。对吧?所以,基本上监控积极、消极和中性情绪随时间的变化,对吧?这对于这种业务非常重要,因为我们想特别看到例如是否有突然增加,是否有负面评论突然增加,如果有,那么我们就需要干预并解决根本原因。对吧?所以,我们将使用的数据集将不是直接的数据集,我们需要编写一个查询。所以,基本上我们需要情绪和餐厅名称,以及评论日期。对吧?所以,第一,情绪在这里。评论时间戳也在这里,因为我们想随时间监控情绪,对吧?但我们没有餐厅名称。我们有餐厅 ID。所以,让我们在此之间进行连接,以及您的 dim restaurant 表。对吧?现在这看起来不错。所以,我们将把 r dot name 作为餐厅名称。并且因为我们想随时间进行,所以必须有一个日期列,这将是评论日期。现在我们想找出每个餐厅随时间的积极、消极和中性评论的数量。对吧?所以,让我们进行 count distinct,case when case when sentiment equals positive then order id else null end,这将是 positive review count。让我们对消极和中性评论也做同样的事情。是的。所以,这将是中性,这将是消极。是的。让我们按 1, 2, 2 进行分组。是的。好的。所以,我们看到每个餐厅在一段时间内,我们现在能够了解随时间的积极、消极和中性评论。所以,让我们继续使用它。我们将创建一个图表。这将是好的。所以,让我将其重命名为评论趋势。是的。这将是您的评论趋势。x 轴将是评论日期,这将是每日。是的,让我们将其命名为日期。y 轴将是 positive review count。是的。让我们将其设为折线图。是的。同样,我们还将添加消极。我们还将添加中性。中性将是黄色。消极将是红色,积极将是绿色。是的。所以,这就是我们看到的,例如。好的,让我快速重命名。积极评论数。这个将是消极评论数,这个将是中性评论数。是的。所以,我们看到例如在这一天,积极评论数为三,而消极和中性评论数为零。是的。下一个图表我们将制作的是每个餐厅的评分分布。所以,例如,一个餐厅获得了多少五星、四星、三星的评分,对吧?所以,让我们先创建数据。我们将使用 gold 表,即您的 D restaurant reviews。所以,这个表基本上有餐厅名称和评分 5、4、3、2、1 的计数。是的。但条形图的工作方式是,对吧?所以,我们想构建它的方式是,在 y 轴上,您将有五星、三、五、四、三、二、一,然后在 x 轴上,您有评论数量,对吧?基本上客户给五星、四星等的次数。对吧?所以,这个表需要,所以,现在,如果我说它已经被透视了,您需要将其取消透视。对吧?这是一个 pandas 特定的术语。您基本上需要将其取消透视。对吧?所以,键或评分标签将是类似五星、四星的东西,然后它将具有计数,客户给五星、四星的次数,对吧?所以,它看起来应该像这样,它看起来应该像餐厅名称、评分标签,然后您有评分计数,对吧?例如,餐厅名称 X,这将是五星,餐厅获得五星的次数是,比如说 10,其他行也是如此,对吧?所以,它应该看起来像这样。现在,如何在 Spark SQL 中进行取消透视操作?是的。有一个叫做 lateral view stack 的东西。是的。它的工作方式基本上是取消透视一定数量的列。对吧?所以,让我们看看它是如何做到的。所以,我们首先将从整个表和 lateral view stack 中选择所有内容。让我们先看看它看起来是什么样的。对吧?所以,我们将说每个主键我想要多少行?对吧?所以,这里的主键是餐厅 ID 本身。每个主键我想要多少行?所以,我想要五行。是的。对于评分 1、2、3、4、5,我想要五行。它们将是什么样的?它们将如下所示。所以,这将是五星。值是什么?值是 rating five count。同样,四星,评分,四计数,三星,依此类推。对吧?我将将其命名为 rating label 和 rating count。所以,这是您的 rating label,这是您的 rating count。让我选择餐厅名称,这将是您的 rating label 和 rating。所以,每个餐厅名称将创建五个 rating label 和五个 rating count。让我们运行它。完美。所以,您看到,让我按餐厅排序,所以您看到每个餐厅创建了五行,即您的 5、4、3、2、1,然后您有您的 rating count,同样适用于所有其他行。是的。所以,让我将其重命名为,让我快速将其重命名为 review count。是的。让我们在这里创建另一个图表。这个图表只是您的评论计数。x 轴,y 轴将是标签,x 轴将是 rating count。是的。让我们更改颜色,也让我们放置标签。所以,这将是 y 轴上的评分本身。这将是您的评论数量。是的。让我们也为每个图表添加一个标题。这将是您的评论情绪趋势,这将是评分分布。是的。所以,让我为特定餐厅选择它,实际上,让我添加评论趋势,这将是按餐厅名称,无论评论计数也应该按餐厅名称过滤。是的。所以,如果我选择另一个,它应该会改变。是的。所以,您看到这个已经改变了。现在我选择另一个,它也应该改变。好的。完美。所以,现在我们已经创建了这个显示餐厅评分分布的图表。接下来,我们将执行问题分类。所以,我们将列出餐厅在用户给出评论时遇到的不同类型的问题。这些可能是您的交付问题、食品质量问题、份量问题或定价问题。是的。所以,首先我们需要创建一个数据集,我们将为此使用 fact reviews,我们之前已经看到我们已经在这里对问题进行了分类,即是交付问题、食品质量、定价或份量问题。是的。所以,让我们将其与
昏暗的餐厅表格,以便我们能够检索餐厅的名称,这将是 R.name,在餐厅名称,然后我们开始对所有问题进行分类,对吧?计数不同的情况,当 fr.fr.issue delivery 等于 true 时,那么它应该是我的订单 ID,它应该被计数,然后是订单 ID,否则为 null n,这将是计数问题交付,对吧?让我们为另一个食物质量、价格和最后一部分大小来做。我们需要对此进行分组,按一分组,让我们看看结果。FR.issue price 无法解决。好的。所以,让我们快速看一下我们拥有的列,问题定价。这应该是定价和问题定价。问题部分大小是正确的。问题食物质量是正确的。所以,让我们运行这个。好的,现在这工作得完美无缺。所以,我们将回到我们的仪表板,让我们在这里放几个图表。这将是您的交付问题。让我们也将其重命名为问题摘要。问题分类而不是问题摘要。是的。所以,我们将在这里使用问题分类。这将是一个简单的计数器。值将是计数问题交付。并且同样,这也将按餐厅名称过滤。所以,我们把它放在这里,让我们放一个标题,上面写着交付问题,问题计数。同样,我们为其他食物质量问题也这样做。这个将是部分大小,部分大小问题计数。然后最后,这个将是定价问题。是的,这看起来不错。现在,让我们来整理一下。完美。我们现在已经将所有指标放置在仪表板上,这些指标对不同问题进行了分类。让我们创建最后两个图表,即最近的正面评论和最近的负面评论。为此,让我们首先创建数据集,我们将使用事实评论。所以,事实评论已经有了我们的评论文本,但它没有餐厅名称。所以,这就是为什么我们将在这里进行连接,我们将简单地放置一个 where,其中 fr.sentiment 等于 positive,并且 fr text 不为 null。所以,这将是 R.name 作为餐厅名称,这将是您的评论文本和评论时间戳。是的。最后,让我们按评论时间戳降序排序。让我们运行这个。好的,这看起来不错。这将是您最近的正面评论。是的,让我们复制这个并创建最近的负面评论。我们只需要将其更改为负面。让我们运行这个。好的,完美。所以,我们将在这里创建两个表。第一个将是最近的正面评论。让我们也将其添加到过滤器列表中。这将按餐厅名称过滤。让我们也为负面评论添加。这也将按餐厅名称过滤。所以,当我点击这个时,我看到,好的,所有这些都可以过滤。这个将是一个简单的表。这次在表中,我们将有评论文本。我们将有评论时间戳,让我们将其命名为评论。是的,让我们将标题设为最近的正面评论。让我们也通过将其着色为绿色来使其更具吸引力。是的,我将复制这个。更改为最近的负面评论。是的,这将是评论。让我们将其更改为深红色。让我们也添加评论时间戳。所以,我们在这里看到评论时间戳,我们也在这里看到评论时间戳。对。所以,现在我们有了由我们的湖仓驱动的两个生产就绪的仪表板。第一个是连锁店绩效洞察。第二个是您的客户评论,对吧?所有这些都建立在我们从头开始创建的数据平台上。让我们快速发布此仪表板。所以,我们将点击发布。然后,您就完成了。所以,这是您的最终仪表板。所以,让我们快速回顾一下我们构建的内容,对吧?两个数据源,Azure 事件中心用于实时流式传输订单,另一个 Azure SQL 数据库用于批量注入。是的。然后您有两种注入方法。第一个是用于流式传输的 Spark 声明式管道,第二个是用于批处理模式的 Lakeflow Connect。是的。然后我们在这里有一个三层奖章架构,即您的铜层用于原始数据,银层用于清理和丰富的事实表,这就是我们实现数据模型的地方,其中一个最有趣的部分是使用 AI 驱动的分析,在事实评论中,我们在 Mosaic AI 上部署了一个模型,并基本上提供一个提示来理解客户一直面临的情绪、问题类型,并将问题分类为,例如,定价、部分大小或交付,对吧?然后我们为聚合业务指标提供了最终的金层,对吧?所有这些都由您的 Unity Catalog 和 Delta Lake 提供支持,用于存储底层数据,对吧?Unity Catalog 具有正确的模式命名约定,用于登陆、铜、银和金。所有这些都使用工作流进行编排,对吧?所以,我们已经看到,我们已经使用了或我们已经使用了工作流来处理流程的这一部分,我们摄取实时订单,然后处理来自铜、银和金的数据。最后,使用 Databricks SQL 来编写我们的查询,进行一些实验,构建我们需要的表,用于 AI BI 仪表板可视化。是的。所以,这个项目的所有代码都在 GitHub 上,以及示例数据集和解释图表架构。链接在描述中,根据您的用例进行调整,并告诉我您从这里构建了什么。对。如果您觉得这个项目能帮助您真实地了解数据工程是什么样的,以及与 Databricks 合作是什么样的,请在 LinkedIn 上分享并标记我。我相信您会喜欢我频道上其他与面试准备、Spark 性能调优、Delta Lake Databricks 相关的视频,对吧?所以,请不要忘记点击喜欢、分享和订阅按钮以及铃铛图标,以获取更多内容。在评论中留下您希望我下次涵盖的主题。