공유 변수 (Shared Variables)

공유 변수 (Shared Variables)

보통 Spark 연산(예: map 또는 reduce)에 전달된 함수가 원격 클러스터 노드에서 실행될 때, 그 함수가 사용하는 모든 변수의 별도 복사본을 갖고 동작해요. 즉 각 머신마다 변수가 복사되고, 원격 머신에서 변수를 갱신해도 그 내용이 드라이버 프로그램으로 다시 전파되지는 않죠. 일반적인 읽기·쓰기 공유 변수를 태스크 전체에 지원하려면 비효율적일 수밖에 없어요.

그래서 Spark는 자주 쓰이는 두 가지 패턴을 위해 제한된 두 종류의 공유 변수를 제공합니다. 바로 브로드캐스트 변수(broadcast variables)누적기(accumulators) 예요.

브로드캐스트 변수 (Broadcast Variables)

브로드캐스트 변수를 사용하면 프로그래머가 읽기 전용 변수를 태스크마다 복사본을 보내는 대신, 각 머신에 캐시해 둘 수 있어요. 예를 들어 큰 입력 데이터셋의 복사본을 모든 노드에 효율적으로 나눠줄 때 쓰면 돼요. Spark는 통신 비용을 줄이기 위해 효율적인 브로드캐스트 알고리즘을 사용해서 변수를 배포하려고 시도합니다.

Spark 액션은 분산 "셔플(shuffle)" 연산으로 구분된 일련의 스테이지들을 거쳐 실행돼요. Spark는 각 스테이지 안의 태스크들이 필요로 하는 공통 데이터를 자동으로 브로드캐스트하고, 이렇게 브로드캐스트된 데이터는 직렬화된 형태로 캐시되었다가 각 태스크 실행 전에 역직렬화됩니다. 즉, 명시적으로 브로드캐스트 변수를 만드는 것은 여러 스테이지의 태스크들이 같은 데이터를 필요로 하거나, 역직렬화된 형태의 데이터 캐시가 중요할 때만 유용해요.

브로드캐스트 변수는 변수 v에서 SparkContext.broadcast(v)를 호출해 만들어요. 브로드캐스트 변수는 v를 감싸는 래퍼이며, 그 값은 value 메서드로 접근할 수 있어요. 아래 코드가 이를 보여줍니다.

Python:

>>> broadcastVar = sc.broadcast([1, 2, 3])
<pyspark.core.broadcast.Broadcast object at 0x102789f10>

>>> broadcastVar.value
[1, 2, 3]

Scala:

scala> val broadcastVar = sc.broadcast(Array(1, 2, 3))
broadcastVar: org.apache.spark.broadcast.Broadcast[Array[Int]] = Broadcast(0)

scala> broadcastVar.value
res0: Array[Int] = Array(1, 2, 3)

Java:

Broadcast<int[]> broadcastVar = sc.broadcast(new int[] {1, 2, 3});

broadcastVar.value();
// returns [1, 2, 3]

브로드캐스트 변수를 만든 뒤에는, 클러스터에서 실행되는 함수 안에서 값 v 대신 이 브로드캐스트 변수를 사용해야 v가 노드들에게 두 번 이상 전달되지 않아요. 또한 브로드캐스트 뒤에는 객체 v를 수정하지 말아야 합니다. 그래야 모든 노드가 브로드캐스트 변수의 같은 값을 보장받을 수 있어요(예: 변수가 나중에 새 노드로 전달되는 경우).

브로드캐스트 변수가 executor에 복사한 리소스를 해제하려면 .unpersist()를 호출해요. 이후에 브로드캐스트를 다시 사용하면 다시 브로드캐스트됩니다. 브로드캐스트 변수가 사용한 모든 리소스를 영구히 해제하려면 .destroy()를 호출하세요. 그 뒤로는 그 브로드캐스트 변수를 사용할 수 없어요. 이 메서드들은 기본적으로 블로킹하지 않는다는 점에 유의하세요. 리소스가 해제될 때까지 기다리려면 호출 시 blocking=true를 지정하면 됩니다.

누적기 (Accumulators)

누적기는 결합 법칙(associative)과 교환 법칙(commutative)을 만족하는 연산으로 "더하기"만 가능한 변수라서 병렬로 효율적으로 지원할 수 있어요. MapReduce에서처럼 카운터나 합계를 구현할 때 씁니다. Spark는 숫자 타입의 누적기를 기본 지원하며, 프로그래머가 새 타입을 추가할 수도 있어요.

사용자는 이름 있는(또는 이름 없는) 누적기를 만들 수 있어요. 아래 그림처럼 이름 있는 누적기(이 예시에서는 counter)는 그 누적기를 수정하는 스테이지의 웹 UI에 표시됩니다. Spark는 각 태스크가 수정한 누적기의 값을 "Tasks" 테이블에 보여줍니다.

UI에서 누적기를 추적하면 실행 중인 스테이지의 진행 상황을 이해하는 데 유용해요 (참고: 이 기능은 Python에서는 아직 지원되지 않습니다).

Python

누적기는 초깃값 v에서 SparkContext.accumulator(v)를 호출해 만들어요. 클러스터에서 실행되는 태스크들은 add 메서드나 += 연산자로 값을 더할 수 있지만, 그 값을 읽을 수는 없어요. 오직 드라이버 프로그램만 value 메서드로 누적기의 값을 읽을 수 있습니다.

아래 코드는 배열의 원소를 합산하는 데 누적기를 사용하는 예시예요.

>>> accum = sc.accumulator(0)
>>> accum
Accumulator<id=0, value=0>

>>> sc.parallelize([1, 2, 3, 4]).foreach(lambda x: accum.add(x))
...
10/09/29 18:41:08 INFO SparkContext: Tasks finished in 0.317106 s

>>> accum.value
10

이 코드는 Int 타입 누적기의 기본 지원을 사용했지만, 프로그래머는 AccumulatorParam을 상속해 자신만의 타입을 만들 수도 있어요. AccumulatorParam 인터페이스에는 두 개의 메서드가 있어요. 데이터 타입의 "영(0) 값"을 제공하는 zero와, 두 값을 더하는 addInPlace입니다. 예를 들어 수학적 벡터를 나타내는 Vector 클래스가 있다면 다음과 같이 작성할 수 있어요.

class VectorAccumulatorParam(AccumulatorParam):
    def zero(self, initialValue):
        return Vector.zeros(initialValue.size)

    def addInPlace(self, v1, v2):
        v1 += v2
        return v1

# Then, create an Accumulator of this type:
vecAccum = sc.accumulator(Vector(...), VectorAccumulatorParam())

Scala / Java

숫자 누적기는 SparkContext.longAccumulator() 또는 SparkContext.doubleAccumulator()를 호출해 각각 Long 또는 Double 타입의 값을 누적할 수 있어요. 클러스터에서 실행되는 태스크들은 add 메서드로 값을 더할 수 있지만, 그 값을 읽을 수는 없어요. 오직 드라이버 프로그램만 value 메서드로 누적기의 값을 읽을 수 있습니다.

아래 코드는 배열의 원소를 합산하는 데 누적기를 사용하는 예시예요.

Scala:

scala> val accum = sc.longAccumulator("My Accumulator")
accum: org.apache.spark.util.LongAccumulator = LongAccumulator(id: 0, name: Some(My Accumulator), value: 0)

scala> sc.parallelize(Array(1, 2, 3, 4)).foreach(x => accum.add(x))
...
10/09/29 18:41:08 INFO SparkContext: Tasks finished in 0.317106 s

scala> accum.value
res2: Long = 10

Java:

LongAccumulator accum = jsc.sc().longAccumulator();

sc.parallelize(Arrays.asList(1, 2, 3, 4)).foreach(x -> accum.add(x));
// ...
// 10/09/29 18:41:08 INFO SparkContext: Tasks finished in 0.317106 s

accum.value();
// returns 10

이 코드는 Long 타입 누적기의 기본 지원을 사용했지만, 프로그래머는 AccumulatorV2를 상속해 자신만의 타입을 만들 수도 있어요. AccumulatorV2 추상 클래스에는 반드시 오버라이드해야 하는 메서드가 몇 가지 있어요. 누적기를 영(0)으로 리셋하는 reset, 누적기에 다른 값을 더하는 add, 같은 타입의 다른 누적기를 합치는 merge입니다. 그 외에 반드시 오버라이드해야 하는 메서드들은 API 문서에 나와 있어요. 예를 들어 수학적 벡터를 나타내는 MyVector 클래스가 있다면 다음과 같이 작성할 수 있어요.

Scala:

class VectorAccumulatorV2 extends AccumulatorV2[MyVector, MyVector] {

  private val myVector: MyVector = MyVector.createZeroVector

  def reset(): Unit = {
    myVector.reset()
  }

  def add(v: MyVector): Unit = {
    myVector.add(v)
  }
  ...
}

// Then, create an Accumulator of this type:
val myVectorAcc = new VectorAccumulatorV2
// Then, register it into spark context:
sc.register(myVectorAcc, "MyVectorAcc1")

Java:

class VectorAccumulatorV2 implements AccumulatorV2<MyVector, MyVector> {

  private MyVector myVector = MyVector.createZeroVector();

  public void reset() {
    myVector.reset();
  }

  public void add(MyVector v) {
    myVector.add(v);
  }
  ...
}

// Then, create an Accumulator of this type:
VectorAccumulatorV2 myVectorAcc = new VectorAccumulatorV2();
// Then, register it into spark context:
jsc.sc().register(myVectorAcc, "MyVectorAcc1");

프로그래머가 자신만의 AccumulatorV2 타입을 정의할 때, 결과 타입이 더해지는 원소들의 타입과 달라질 수 있다는 점에 유의하세요.

경고: Spark 태스크가 끝나면 Spark는 그 태스크의 누적 갱신을 누적기에 병합하려고 해요. 만약 그 병합이 실패하면 Spark는 실패를 무시하고 여전히 태스크를 성공으로 표시한 뒤 나머지 태스크를 계속 실행합니다. 따라서 버그가 있는 누적기는 Spark 잡에 영향을 주지는 않지만, 잡이 성공했음에도 누적기가 올바르게 갱신되지 않을 수 있어요.

액션 안에서 수행되는 누적기 갱신에 대해서만, Spark는 각 태스크의 갱신이 정확히 한 번만 적용된다는 것을 보장해요. 즉 재시작된 태스크는 값을 갱신하지 않습니다. 변환(transformation)에서 사용자는 태스크나 잡 스테이지가 다시 실행되면 각 태스크의 갱신이 두 번 이상 적용될 수 있다는 점을 알아야 해요.

누적기는 Spark의 지연 평가(lazy evaluation) 모델을 바꾸지 않아요. RDD에 대한 연산 안에서 누적기를 갱신한다면, 그 값은 해당 RDD가 액션의 일부로 계산될 때만 갱신됩니다. 따라서 map() 같은 지연 변환 안에서 누적기를 갱신한다 해도 그 갱신이 실행되리라는 보장은 없어요. 아래 코드 조각이 이 특성을 보여줍니다.

Python:

accum = sc.accumulator(0)
def g(x):
    accum.add(x)
    return f(x)
data.map(g)
# Here, accum is still 0 because no actions have caused the `map` to be computed.

Scala:

val accum = sc.longAccumulator
data.map { x => accum.add(x); x }
// Here, accum is still 0 because no actions have caused the map operation to be computed.

Java:

LongAccumulator accum = jsc.sc().longAccumulator();
data.map(x -> { accum.add(x); return f(x); });
// Here, accum is still 0 because no actions have caused the `map` to be computed.