本文属于机器翻译版本。若本译文内容与英语原文存在差异,则一律以英文原文为准。
使用 PySpark 分析模板在配置的表上运行 PySpark 作业
此过程演示如何使用 AWS Clean Rooms 控制台中的 PySpark 分析模板使用自定义分析规则分析配置的表。
使用 PySpark 分析模板在配置的表上运行 PySpark 作业
登录 AWS 管理控制台 并打开 AWS Clean Rooms 控制台,地址为https://console.aws.amazon.com/cleanrooms
-
在左侧导航窗格中,选择协作。
-
选择您的成员能力状态为 “运行作业” 的协作。
-
在分析选项卡的表格部分下,查看表及其关联的分析规则类型(自定义分析规则)。
-
在分析部分下,对于分析模式,选择运行分析模板。
-
从 “ PySpark 分析模板” 下拉列表中选择分析模板。
PySpark 分析模板中的参数将自动填充到定义中。
-
如果分析模板定义了参数,请在参数下提供参数值:
-
对于每个参数,查看参数名称和默认值(如果已配置)。
-
为要覆盖的每个参数输入一个值。
注意
如果您未提供值但存在默认值,则将使用默认值。
重要
参数值最多可为 1,000 个字符,并支持 UTF-8 编码。所有参数值都被视为字符串,并通过上下文对象传递给您的用户脚本。
确保您的用户脚本安全地验证和处理参数值。有关安全参数处理的更多信息,请参阅使用 PySpark 分析模板中的参数。
-
-
指定支持的工作人员类型和工作人员数量。
使用下表来确定用例所需的人员类型和人数。
Worker 类型 vCPU 内存(GB) 存储(GB) 工作线程数 洁净室处理单元总数 (CRPU) CR.1X(默认值) 4 30 100 4 8 128 256 CR.4X 16 120 400 4 32 32 256 注意
不同的工人类型和工人数量有相关的成本。要了解有关定价的更多信息,请参阅AWS Clean Rooms 定价
。 -
指定支持的 Spark 属性。
-
选择添加 Spark 属性。
-
在 Spark 属性对话框中,从下拉列表中选择一个属性名称并输入一个值。
下表提供了每个属性的定义。
有关 Spark 属性的更多信息,请参阅 Apache Spark 文档
中的 Spark 属性。 注意
您最多可以配置 50 个 Spark 属性。每个属性值最多可为 500 个字符。
属性名称 说明 默认值 spark.task.max 失败
控制任务在任务失败之前可以连续失败多少次。需要一个大于或等于 1 的值。允许的重试次数等于该值减去 1。如果任何尝试成功,失败次数将重置。不同任务的失败次数不会累积到此上限。
4
spark.sql.files.max PartitionBytes
设置从 Parquet、JSON 和 ORC 等基于文件的源读取时打包到单个分区的最大字节数。
128MB
spark.hadoop.fs.s3.max 重试次数
设置 Amazon S3 文件操作的最大重试次数。
(无)
火花网络超时
为所有网络交互设置默认超时。如果未配置以下超时设置,则会覆盖这些设置:
-
Spark.storage.block ManagerHeartbeatTimeoutMs
-
spark.shuffle.io.连接超时
-
Spark.rpc.askTimeout
-
spark.rpc.lookupTime
120 秒
spark.rdd.compress
指定是否使用 spark.io.compression.codec 压缩序列化的 RDD 分区。适用于 Java 和 Scala 中的 StorageLevel.MEMORY _ONLY_SER,或 StorageLevel.MEMORY Python 中的 _ONLY。减少了存储空间,但需要额外的 CPU 处理时间。
false
spark.shuffle.spill.compress
指定是否使用 spark.io.compression.codec 压缩随机泄漏数据。
true
spark.shuffle.compress
指定是否压缩地图输出文件。压缩使用 spark.io.compression.codec。
true
spark.shuffle.service.index.cache.size
设置缓存大小限制,除非另行指定,否则以字节为单位。
100 m
spark.shuffle.io.max重试次数
设置因异常而失败的读取的最大重试次数。 IO-related
3
spark.shuffle.io.retryWait
设置两次重试提取之间的等待时间。默认情况下,重试造成的最大延迟为 15 秒,计算方法为 maxRetries * retryWait。
5 秒
spark.shuffle.io.连接超时
如果仍有未完成的提取请求但频道上没有流量,则将随机服务器和客户端之间建立的连接的超时时间设置为空闲和关闭。
(spark.network.timeout 的值
Spark.driver.max ResultSize
设置每个 Spark 操作的所有分区序列化结果的总大小限制,以字节为单位。应至少为 1M,或者 0 表示无限制。
1g
火花。记忆。分数
设置用于执行和存储的(堆空间-300MB)的比例。该值越低,发生泄漏和缓存数据驱逐的频率越高。建议将其保留为默认值。
0.6
Spark.scheduler.mode
设置提交给该任务的任务之间的调度模式 SparkContext。可以设置为 FAIR 以使用公平共享,而不是依次排队作业。支持的值:FAIR、FIFO。
FIFO
spark.sql.adaptive.adv PartitionSizeInBytes
当 spark.sql.adaptive.enabled 为真时,设置自适应优化期间随机分区的目标大小(以字节为单位)。在合并小分区或拆分倾斜分区时控制分区大小。
(spark.sql.adaptive.shuffle.target PostShuffleInputSize 的值)
spark.sql.aptive.auto BroadcastJoinThreshold
设置连接期间向工作节点广播的最大表大小(以字节为单位)。仅适用于自适应框架。使用与 spark.sql. BroadcastJoinThreshold auto 相同的默认值。设置为 -1 以禁用广播。
(无)
spark.sql.adaptive.coalesce Partitions.enabled
指定是否根据 spark.sql.adaptive.advisory 合并连续的洗牌分区以优化任务大小。PartitionSizeInBytes 需要 spark.sql.adaptive.enabled 为真。
true
spark.sql.adaptive.coalesce Partitions.initialPartitionNum
定义合并前洗牌分区的初始数量。要求 spark.sql.adaptive.enabled 和 spark.sql.adaptive.coalesce 都为真。Partitions.enabled 默认为 spark.sql.shuffle.partitions 的值。
(无)
spark.sql.adaptive.coalesce Partitions.minPartitionSize
设置合并随机播放分区的最小大小,以防止分区在自适应优化期间变得太小。
1 MB
spark.sql.adaptive.coalesce Partitions.parallelismFirst
指定在分区合并期间是否根据集群并行度而不是 spark.sql.adaptive. PartitionSizeInBytes advisory 来计算分区大小。生成的分区大小小于配置的目标大小,以最大限度地提高并行度。我们建议在繁忙的群集上将其设置为 false,通过防止过多的小任务来提高资源利用率。
true
spark.sql.adaptive.enab
根据准确的运行时统计信息,指定是否启用自适应查询执行以在查询执行期间重新优化查询计划。
true
spark.sql.aptive.force OptimizeSkewedJoin
指定 OptimizeSkewedJoin 即使引入了额外的随机播放也要强制启用。
false
spark.sql.aptive.local ShuffleReader.enabled
指定在不需要随机分区时(例如从排序合并联接转换为广播哈希联接之后)是否使用本地随机播放阅读器。需要 spark.sql.adaptive.enabled 为真。
true
spark.sql.adaptive.max ShuffledHashJoinLocalMapThreshold
设置用于构建本地哈希映射的最大分区大小(以字节为单位)。在以下情况下,将随机哈希连接优先于排序合并联接:
-
此值等于或超过 spark.sql.adaptive.advisory PartitionSizeInBytes
-
所有分区大小都在此限制之内
覆盖 spark.sql.join.prefer 设置SortMergeJoin 。
0 字节
spark.sql.adaptive.优化 SkewsInRebalancePartitions.enabled
指定是否通过基于 spark.sql.adaptive.advisory 将倾斜的随机分区拆分成较小的分区来优化这些分区。PartitionSizeInBytes需要 spark.sql.adaptive.enabled 为真。
true
spark.sql.aptive.rebalance PartitionsSmallPartitionFactor
定义拆分期间合并分区的大小阈值系数。小于此系数乘以 spark.sql.adaptive. PartitionSizeInBytes advisory 的分区将被合并。
0.2
spark.sql.adaptive.skew Join.enabled
指定是否通过拆分和可选地复制倾斜分区来处理随机连接中的数据倾斜。适用于排序合并和随机哈希连接。需要 spark.sql.adaptive.enabled 为真。
true
spark.sql.adaptive.skew Join.skewedPartitionFactor
确定决定分区倾斜的大小因子。当分区的大小超过这两个分区时,分区就会倾斜:
-
该因子乘以分区大小中位数
-
spark.sql.adaptive.skew 的值 Join.skewedPartitionThresholdInBytes
5
spark.sql.adaptive.skew Join.skewedPartitionThresholdInBytes
为识别倾斜分区设置大小阈值(以字节为单位)。当分区的大小超过这两个分区时,分区就会倾斜:
-
这个阈值
-
分区大小中位数乘以 spark.sql.adaptive.skew Join.skewedPartitionFactor
我们建议将此值设置为大于 spark.sql.adaptive.advisory PartitionSizeInBytes。
256MB
spark.sql.广播超时
控制广播加入期间广播操作的超时时间(以秒为单位)。
300 秒
spark.sql.cbo.enabled
指定是否启用基于成本的优化 (CBO) 以进行计划统计估计。
false
spark.sql.cbo.join Reorder.dp.star.filter
指定在基于成本的联接枚举期间是否应用星型联接过滤器启发式方法。
false
spark.sql.cbo.join Reorder.dp.threshold
设置动态规划算法中允许的最大连接节点数。
12
spark.sql.cbo.join Reorder.enabled
指定是否在基于成本的优化 (CBO) 中启用联接重新排序。
false
spark.sql.cbo.plan Stats.enabled
指定在逻辑计划生成期间是否从目录中提取行数和列统计信息。
false
spark.sql.cbo.star SchemaDetection
指定是否启用基于星形架构检测的联接重新排序。
false
spark.sql.files.max PartitionNum
为基于文件的源(Parquet、JSON 和 ORC)设置分区的最大目标分区数。当初始计数超过此值时,重新缩放分区。这是建议的目标,不是保证的上限。
(无)
spark.sql.files.max RecordsPerFile
设置写入单个文件的最大记录数。当设置为零或负值时,没有限制。
0
spark.sql.files.min PartitionNum
为基于文件的源(Parquet、JSON 和 ORC)设置拆分文件分区的目标最小数量。默认为 spark.sql.le NodeDefaultParallelism af。这是建议的目标,不是保证的上限。
(无)
spark.sql.in MemoryColumnarStorage.batchSize
控制列式缓存的批次大小。增加大小可以提高内存利用率和压缩率,但会增加内存不足错误的风险。
10000
spark.sql.in MemoryColumnarStorage.compressed
指定是否根据数据统计信息自动为列选择压缩编解码器。
true
spark.sql.in MemoryColumnarStorage.enableVectorizedReader
指定是否为列式缓存启用矢量化读取。
true
spark.sql.legacy.alLOW HashOnMapType
指定是否允许对地图类型数据结构进行哈希运算。此传统设置保持了与旧 Spark 版本的地图类型处理的兼容性。
(无)
spark.sql.legacy.alLOW NegativeScaleOfDecimal
指定是否允许在十进制类型定义中使用负比例值。此传统设置保持了与支持负十进制刻度的旧 Spark 版本的兼容性。
(无)
spark.sql.legacy.cast ComplexTypesToString.enabled
指定是否启用将复杂类型转换为字符串的传统行为。保持与旧 Spark 版本的类型转换规则的兼容性。
(无)
spark.sql.legacy.char VarcharAsString
指定是否将 CHAR 和 VARCHAR 类型视为字符串类型。此传统设置可与旧版 Spark 版本的字符串类型处理兼容。
(无)
spark.SQL.legacy.create EmptyCollectionUsingStringType
指定是否使用字符串类型元素创建空集合。此传统设置保持了与旧 Spark 版本的集合初始化行为的兼容性。
(无)
spark.sql.legacy.exponent LiteralAsDecimal.enabled
指定是否将指数文字解释为十进制类型。此传统设置保持了与旧版 Spark 版本的数字文字处理的兼容性。
(无)
spark.sql.legacy.json.allow EmptyString.enabled
指定是否允许在 JSON 处理中使用空字符串。此传统设置保持了与旧版 Spark 版本的 JSON 解析行为的兼容性。
(无)
spark.sql.legacy.parquet.int96 RebaseModeInRead
指定在读取 Parquet 文件时是否使用传统的 INT96 时间戳变基模式。此传统设置保持了与旧 Spark 版本的时间戳处理的兼容性。
(无)
spark.sql.legacy.time ParserPolicy
控制时间解析行为以实现向后兼容。这个传统设置决定了如何从字符串中解析时间戳和日期。
(无)
spark.SQL.legacy.type Coercion.datetimeToString.enabled
指定在将日期时间值转换为字符串时是否启用传统类型强制行为。保持与旧版 Spark 版本的日期时间转换规则的兼容性。
(无)
spark.sql.max SinglePartitionBytes
以字节为单位设置最大分区大小。规划器为较大的分区引入了洗牌操作以提高并行性。
128 米
spark.sql.MetadataCachettlSecon
控制元数据缓存的生存时间 (TTL)。适用于分区文件元数据和会话目录缓存。需要:
-
大于零的正值
-
spark.sql.Catalog实现设置为 hive
-
spark.sql.hive. PartitionFileCacheSize filesource 大于零
-
spark.sql.hive.manage 设置为 true FilesourcePartitions
-1000 毫秒
spark.sql 优化器折叠 ProjectAlwaysInline
指定是否折叠相邻的投影和行内表达式,即使它会导致重复。
false
spark.sql.optimizer.dyn PartitionPruning.enabled
指定是否为用作联接键的分区列生成谓词。
true
spark.sql.optimizer.enable CsvExpressionOptimization
指定是否通过从 from_csv 操作中删除不必要的列来优化 SQL 优化器中的 CSV 表达式。
true
spark.sql.optimizer.enable JsonExpressionOptimization
通过以下方式指定是否在 SQL 优化器中优化 JSON 表达式:
-
从 from_json 操作中删除不必要的列
-
简化 from_json 和 to_json 的组合
-
优化 named_struct 操作
true
spark.sql.optimizer.excludedRules
定义要禁用的优化器规则,由逗号分隔的规则名称标识。某些规则无法禁用,因为它们是正确性所必需的。优化器会记录哪些规则已成功禁用。
(无)
spark.SQL.optimizer.runtime.bloom Filter.applicationSideScanSizeThreshold
设置在应用程序端注入 Bloom 过滤器所需的最小聚合扫描大小(以字节为单位)。
10GB
spark.SQL.optimizer.runtime.bloom Filter.creationSideThreshold
定义在创建端注入 Bloom 滤镜的最大大小阈值。
10MB
spark.SQL.optimizer.runtime.bloom Filter.enabled
指定当随机连接的一端具有选择性谓词时,是否插入 Bloom 过滤器以减少随机播放数据。
true
spark.SQL.optimizer.runtime.bloom Filter.expectedNumItems
定义运行时 Bloom 过滤器中预期项目的默认数量。
1000000
spark.SQL.optimizer.runtime.bloom Filter.maxNumBits
设置运行时 Bloom 过滤器允许的最大位数。
67108864
spark.SQL.optimizer.runtime.bloom Filter.maxNumItems
设置运行时 Bloom 过滤器中允许的最大预期项目数。
4000000
spark.SQL.optimizer.runtime.bloom Filter.numBits
定义运行时 Bloom 过滤器中使用的默认位数。
8388608
spark.sql.optimizer.runtime.row LevelOperationGroupFilter.enabled
指定是否为行级操作启用运行时组筛选。允许数据源:
-
使用数据源筛选器删除整组数据(例如文件或分区)
-
执行运行时查询以识别匹配的记录
-
丢弃不必要的群组以避免昂贵的重写
限制:
-
并非所有表达式都能转换为数据源筛选器
-
某些表达式需要 Spark 评估(例如子查询)
true
spark.SQL.optimizer.runtim Filter.number.threshold
设置注入的运行时过滤器的总数(非 DPP)。这是为了防止驱动程序 OOM 使用过多 Bloom 过滤器。
10
spark.SQL.optimizer.runtim Filter.semiJoinReduction.enabled
指定当随机连接的一端具有选择性谓词时,是否插入半连接以减少随机播放数据。
false
spark.sql.parquet.aggregatePush
指定是否将聚合向下推送到 Parquet 进行优化。支持:
-
布尔值、整数、浮点数和日期类型的最小值和最大值
-
所有数据类型的 COUNT
如果任何 Parquet 文件页脚中缺少统计信息,则会引发异常。
false
spark.sql.parquet.columnar ReaderBatchSize
控制每个 Parquet 矢量化读取器批次中的行数。选择一个平衡性能开销和内存使用量的值,以防止出现内存不足错误。
4096
spark.sql.parquet.enable VectorizedReader
指定是否启用矢量化 Parquet 解码。
true
spark.sql.shuffle.分区
设置联接或聚合期间数据洗牌的默认分区数。无法在结构化流式查询从同一检查点位置重新启动之间进行修改。
200
spark.sql.shuffled HashJoinFactor
定义用于确定 shuffle 哈希加入资格的乘法系数。当小边数据大小乘以该系数小于大边数据大小时,将选择随机哈希连接。
3
spark.sql.sources.par PartitionDiscovery.threshold
使用基于文件的源(Parquet、JSON 和 ORC)设置驱动程序端文件列表的最大路径数。如果在分区发现期间超过该值,则使用单独的 Spark 分布式作业列出文件。
32
spark.sql.statistics.histics.en
指定在列统计计算期间是否生成等高直方图以提高估计精度。除了基本列统计数据所需的表扫描外,还需要进行额外的表扫描。
false
火花。动态 Allocation.executorIdleTimeout
设置启用动态分配后,执行器在移除之前必须处于空闲状态的持续时间。
60s
火花。动态 Allocation.schedulerBacklogTimeout
设置启用动态分配时在请求新执行者之前必须积压待处理任务的持续时间。
1s
火花。动态 Allocation.sustainedSchedulerBacklogTimeout
与 spark.dynamic 相同Allocation.schedulerBacklogTimeout,但仅用于后续的执行器请求。
(spark.dynamic Allocation.schedulerBacklogTimeout 的值)
spark.scheduler.min RegisteredResourcesRatio
设置计划开始之前等待的注册资源(注册资源/预期资源总数)的最小比率。指定为介于 0.0 和 1.0 之间的双精度。无论是否达到最低资源比例,调度开始前等待的最大时间都由 spark.scheduler. RegisteredResourcesWaitingTime max 控制。
0.8
spark.scheduler.max RegisteredResourcesWaitingTime
设置在开始调度之前等待资源注册的最大时间。
30 秒
spark.sql.hive.metastore PartitionPruningFallbackOnException
指定是否回退到从 Hive 元数据仓获取所有分区,并在从元数据仓遇到 MetaException 分区时在 Spark 客户端执行分区修剪。
false
spark.sql.cro Join.enabled
指定是否允许包含笛卡尔积且没有明确的 CROSS JOIN 语法的查询。
true
spark.sql.analyzer.maxIteration
设置查询分析器在放弃之前运行的最大迭代次数。较高的值允许分析器处理非常大的或深度嵌套的查询。
100
spark.sql.dataprefetch.filescan.max ParallelismPerTask
设置扫描文件时每项任务可同时预取的最大文件拆分次数。
4
spark.sql.iceberg.data-prefetch.已启用
指定在读取Iceberg表时是否启用数据预取优化。
true
spark.sql.legacy.null ValueWrittenAsQuotedEmptyStringCsv
指定是否恢复在 CSV 输出中将空值写入带引号的空字符串的传统行为。当为 false 时,Spark 将空值写入未加引号的空字符串。
false
spark.max RemoteBlockSizeFetchToMem
设置大小阈值,超过该阈值,Spark 会将远程块提取到磁盘而不是内存。这样可以避免单个大请求消耗过多内存。
200 m
spark.emr-serverless.allacion.batch.size
设置在每轮执行器分配中同时请求的执行者数量。
20
属性名称 说明 默认值 spark.SQL.auto BroadcastJoinThreshold
设置连接期间向工作节点广播的最大表大小(以字节为单位)。设置为 -1 以禁用广播。
10MB
spark.io.compression.codecs
设置用于压缩 RDD 分区、事件日志、广播变量和随机播放输出等内部数据的编解码器。支持的值:lz4、snappy、zstd、gzip。
lz4
spark.sql.Session.Time
定义会话时区,用于处理字符串文字和 Java 对象转换中的时间戳。接受:
-
Region-based area/city 格式化的 ID(例如 America/Los _Angeles)
-
以 (+/-) HH、(+/-) 或 (+/-) HH:mm HH:mm:ss 格式的区域偏移量(例如 -08 或 + 01:00)
-
UTC 或 Z 作为 + 00:00 的别名
(本地时区的值)
spark.cleanrooms.executor.mem OverheadFactor
设置用于确定 spark.executor.memory 和 spark.executor.Memory 和 spark.executor.MemoryOveeder.MemoryOveen 指定为介于 0.0 和小于 1.0 之间的双精度。
0.1
spark.cleanrooms.驱动程序.内存 OverheadFactor
设置用于确定 spark.driver.memory 和 spark.driver.Memory 和 Spark.driver.MemoryOveer.Oveer 指定为介于 0.0 和小于 1.0 之间的双精度。
0.1
Spark.memory.storageFr
设置免受驱逐影响的存储内存量,表示为 spark.memory.fraction 预留的区域大小的一小部分。该值越高,可用于执行的工作内存就越少,任务溢出到磁盘的频率可能更高。建议将其保留为默认值。
0.5
Spark.rpc.askTimeout
设置 RPC 任务操作在超时之前等待的持续时间。
(Spark.network.timeout 的值)
spark.executor.HeartbeatIn
设置每个执行者向驱动程序发送心跳之间的间隔。心跳让驱动程序知道执行器还活着,并使用正在进行的任务的指标对其进行更新。spark.executor.HeartbeatInterval 应该明显小于 spark.network.timeout。
10s
spark.stage.max ConsecutiveAttempts
设置一个阶段中止之前允许的连续阶段尝试次数。
4
spark.task.cpus
设置为每项任务分配的内核数。
1
spark.shuffle 文件缓冲区
除非另行指定,否则以 KiB 为单位设置每个随机文件输出流的内存缓冲区的大小。这些缓冲区减少了创建中间 shuffle 文件时进行的磁盘搜索和系统调用的次数。
32k
spark.reducer.max SizeInFlight
设置从每个 reduce 任务中同时提取的地图输出的最大大小,除非另有指定,否则以 MiB 为单位。由于每个输出都需要一个缓冲区才能接收,因此这表示每个 reduce 任务的固定内存开销,因此除非您有大量内存,否则请保持较小的内存开销。
48 米
-
-
(可选)对于计算付款人,选择支付任务计算费用的协作成员。
注意
如果协作中只有一个职位计算付款人候选人,则默认为该付款人。
-
选择运行。
注意
如果可以接收结果的成员未配置作业结果设置,则无法运行作业。
-
继续调整参数并再次运行作业,或者选择 + 按钮在新选项卡中启动新作业。