小菜鸟

java菜鸟号正在起航

Scala 泛型:类型安全与灵活性的平衡

泛型是 Scala 中实现类型抽象的核心机制,它允许在定义类、特质、方法时使用类型参数,从而编写与具体类型无关的通用代码。Scala 泛型不仅支持类似 Java 的上下界约束,还引入了协变、逆变等独特特性,进一步增强了类型系统的灵活性。本文将全面解析 Scala 泛型的用法与高级特性。

泛型基础:类型参数的定义与使用

Scala 泛型使用方括号 [] 定义类型参数,可用于类、特质、方法等,实现代码的复用与类型安全。

泛型类

在类定义中声明类型参数,使类能处理多种类型的数据:

// 定义泛型类 Box,类型参数为 T
class Box[T](val content: T) {
  // 方法返回值使用类型参数 T
  def getContent: T = content
}

// 创建不同类型的 Box 实例
val intBox = new Box[Int](10)
val strBox = new Box[String]("hello")

println(intBox.getContent)  // 输出:10(类型为 Int)
println(strBox.getContent)  // 输出:hello(类型为 String)

泛型方法

在方法中声明类型参数,使方法能独立于类的类型参数处理不同类型:

object GenericUtils {
  // 泛型方法:交换数组中两个位置的元素
  def swap[T](array: Array[T], i: Int, j: Int): Unit = {
    val temp = array(i)
    array(i) = array(j)
    array(j) = temp
  }
}

// 测试泛型方法
val intArray = Array(1, 2, 3)
GenericUtils.swap(intArray, 0, 2)
println(intArray.mkString(","))  // 输出:3,2,1

val strArray = Array("a", "b", "c")
GenericUtils.swap(strArray, 1, 2)
println(strArray.mkString(","))  // 输出:a,c,b
阅读全文 »

J2EE 分布式事务规范:JTA、JTS 与两阶段提交协议

在分布式系统中,保证跨多个资源(如数据库、消息队列)的事务一致性是核心挑战。J2EE 提供了两套规范支持分布式事务:JTA(Java Transaction API)JTS(Java Transaction Service)。其中,JTA 定义了分布式事务的高层接口,JTS 则规定了事务管理器的实现标准,而两阶段提交协议(2PC) 是实现分布式事务强一致性的核心机制。本文将详细解析这些概念及其工作原理。

JTA 与 JTS:分布式事务的规范体系

J2EE 的分布式事务规范由 JTA 和 JTS 共同构成,二者分工明确,形成 “接口定义 - 实现规范” 的层级关系。

JTA:分布式事务的高层 API

JTA(javax.transaction 包)是一套与具体实现无关、与协议无关的高层接口规范,定义了分布式事务的核心操作(如开始、提交、回滚),供开发者与事务管理器交互。其核心目标是屏蔽底层资源(如数据库、消息队列)的差异,提供统一的事务编程模型。

JTA 主要包含以下接口:

接口 作用描述
UserTransaction 供应用程序直接使用的事务操作接口,提供 begin()commit()rollback() 等方法,控制事务生命周期。
Status 定义事务的状态常量(如 STATUS_ACTIVESTATUS_COMMITTEDSTATUS_ROLLEDBACK 等),用于表示事务当前状态。
Synchronization 允许应用程序在事务提交或回滚前后注册回调逻辑(beforeCompletion()afterCompletion()),用于资源清理或日志记录。
Transaction 代表一个分布式事务对象,提供获取事务状态、注册同步器等方法,主要由事务管理器内部使用。
TransactionManager 由事务管理器实现,负责事务的创建、传播和管理,应用程序通常不直接使用,而是通过 UserTransaction 间接交互。
UserTransaction 核心方法示例:
阅读全文 »

Spark 广播变量:大字典表别让每个 Task 都传一份,发给每个 Executor 一份就够了

Spark 里每个 Task 执行时,闭包里的变量会被复制一份到该 Task。如果变量很小(几 KB),没问题。但如果变量很大(比如 1GB 的字典表),100 个 Task 就是 100 份拷贝。

分发方式 拷贝数 网络传输 内存占用
普通闭包 每个 Task 一份 大(100 × 1GB) 大(100 × 1GB)
广播变量 每个 Executor 一份 小(Executor 数 × 1GB) 小(Executor 数 × 1GB)

广播变量就是:把大对象从 Driver 发给每个 Executor,只发一次,该 Executor 上的所有 Task 共享这一份。

广播变量的基本用法

// 1. 创建广播变量
val dict = Map("a" -> 10, "b" -> 20, "c" -> 30)
val bcDict = sc.broadcast(dict)

// 2. Task 里用 broadcast.value 取数据
val rdd = sc.parallelize(List("a", "b", "c", "d"))
val result = rdd.map(key => bcDict.value.getOrElse(key, 0))

// 3. 不需要时释放(可选)
bcDict.unpersist()

关键点:

  • 广播变量是只读的——Executor 上只能读,不能改
  • 广播变量的值在 Executor 上只加载一次,所有 Task 共享
  • 使用 .value 访问广播的内容

广播变量 vs 普通闭包

//  普通闭包:每个 Task 都传一份 dict
val dict = Map("a" -> 10, "b" -> 20)   // 假设 1GB
val rdd = sc.parallelize(1 to 1000, 100)  // 100 个 Task
rdd.map(key => dict.getOrElse(key, 0)).collect()
// 网络传输:100 × 1GB = 100GB

// 广播变量:每个 Executor 只传一份
val bcDict = sc.broadcast(dict)
rdd.map(key => bcDict.value.getOrElse(key, 0)).collect()
// 网络传输:Executor 数 × 1GB(假设 10 个 Executor = 10GB)

典型使用场景

1. Map 端 Join(小表广播)

小表 Join 大表时,把小表广播到每个 Executor,在 Map 端完成关联,不走 Shuffle。

阅读全文 »

Spark 累加器:Driver 的变量在 Executor 上改了没用?用累加器把数”收”回来

在 Spark 里,你可能会写出这样的代码:

var sum = 0
rdd.foreach(x => sum += x)
println(sum)   // 输出:0

直觉上应该是所有元素的和,但结果是 0。因为 sum 在 Driver 端定义,被复制到每个 Executor 上各自修改,改的是副本,Driver 端的 sum 纹丝不动。

累加器就是解决这个问题的:Driver 发一个”计数器”给 Executor,Executor 往里加数,Driver 最后把总数收回来。

累加器就干一件事:把 Executor 上的值聚合回 Driver。

累加器的基本用法

1. 创建累加器

val sumAcc = sc.longAccumulator("Sum")   // 长整型累加器
val countAcc = sc.doubleAccumulator("Count")  // 双精度累加器
val errAcc = sc.collectionAccumulator[String]("Errors")  // 集合累加器

2. Executor 上累加

rdd.foreach(x => sumAcc.add(x))   // 每个 Task 往累加器里加数

3. Driver 读取结果

println(sumAcc.value)   // 所有 Task 累加的总和

完整示例:

val rdd = sc.parallelize(1 to 100)
val sumAcc = sc.longAccumulator("Sum")

rdd.foreach(x => sumAcc.add(x))

println(sumAcc.value)   // 5050

累加器的类型

Spark 内置三种累加器:

类型 创建方法 用途
长整型 sc.longAccumulator("name") 计数、求和(int/long)
双精度 sc.doubleAccumulator("name") 求和(float/double)
集合 sc.collectionAccumulator[T]("name") 收集错误信息、特殊值

集合累加器最实用: 收集处理过程中遇到的异常数据。

阅读全文 »

Scala 闭包与函数柯里化:函数式编程的高级特性

闭包(Closure)和函数柯里化(Currying)是 Scala 函数式编程中的两个核心概念。它们不仅增强了函数的灵活性,还为代码复用和模块化提供了强大支持。本文将深入解析闭包的本质、函数柯里化的实现及其应用场景。

闭包(Closure):函数与环境的结合体

闭包是指一个函数与其引用的外部变量形成的整体。即使外部变量脱离了原有的作用域,只要闭包存在,这些变量就会被保留并供函数使用。

闭包的定义与示例

// 定义一个返回函数的高阶函数
def makeAdder(base: Int): Int => Int = {
  // 匿名函数引用了外部变量 base
  (x: Int) => base + x
}

// 创建闭包:base = 10 被保留在闭包中
val add10 = makeAdder(10)

// 调用闭包:使用保留的 base = 10
println(add10(5))  // 输出:15(10 + 5)
println(add10(3))  // 输出:13(10 + 3)

// 另一个闭包:base = 20 被保留
val add20 = makeAdder(20)
println(add20(5))  // 输出:25(20 + 5)

核心原理

  • makeAdder 是一个高阶函数,返回一个匿名函数。
  • 匿名函数 (x: Int) => base + x 引用了外部参数 base
  • makeAdder(10) 被调用时,base 被固定为 10,与匿名函数绑定形成闭包 add10
  • 即使 makeAdder 执行结束,add10 仍能访问 base = 10,因为闭包保留了对该变量的引用。

闭包的本质

闭包本质上是一个携带状态的函数对象。在 Scala 中,闭包会被编译为 FunctionN 特质的实现类(如 Function1Function2),外部变量会作为该对象的字段被保存。

阅读全文 »
0%