데보션앱 소개페이지 바로가기
로그인 선택

신고하기

CLOSE
신고사유 (대표 사유 1개)
상세내용 (선택)
0/200
  • 신고한 게시글은 더 이상 보이지 않습니다.
  • 이용약관과 운영정책에 따라 신고사유에 해당하는지 검토 후 조치됩니다.
  • 허위 신고인 경우, 신고자의 서비스 이용이 제한될 수 있으니 유의하시어 신중하게 신고해 주세요.
(이 회원이 작성한 모든 댓글과 커뮤니티 게시물이 보이지 않고, 알림도 오지 않습니다.)

미리보기

커뮤니티

      1,234

      badge 23.06.15

      글 등록

      카테고리를 선택해주세요.

      DEVOTEE를 활성화 시키면
      지금 작성한 커뮤니티 글에 대해 1개의 댓글을 달아줍니다.

      버튼을 누르면 글 수정 시 ChatGPT가 작성한 댓글이 수정됩니다.

      임시저장함에 저장되었습니다. 저장일시 : 2022.5.17 14:29:08

      임시저장함

      제목을 선택하시면 이어서 작성이 가능하며,
      최대 20건까지 저장합니다.
      컨텐츠 유형, 제목, 저장일시, 삭제로 이뤄진 임시저장 목록
      컨텐츠 유형 제목 저장일 삭제

      데보션 블로그 게재 요청

      CLOSE
      • *
      • *

      본인인증

      효율적인 데보션 서비스 이용 및
      고객님의 소중한 개인정보보호를 위해
      본인인증을 진행해주세요. 본인인증 미 진행 시 로그인이 제한됩니다.
      본인인증 실패

      본인인증 로그인에 실패하였습니다.
      회원이 아니시거나 본인인증 등록이
      완료되지 않은 사용자입니다.

      회원정보 연결

      Spark 3.3.0 review

      jungbbong 22.08.16
      1,773 12 0

      *해당 글은 spark 3.3.0 릴리즈노트의 1608개의 JIRA에서 개선된 내용을 요약 정리 및 Data+AI Summit 2022에서 발표한 내용들을 합쳐서 만든 글입니다.


      최근 spark은 python사용자들이 늘어남에 따라서 python에 대한 지원을 좀 더 확장하도록 프로젝트가 신설(Project Zen)되어서 진행중이다. 해당 프로젝트의 결과로 이번에 릴리즈된 spark 3.3.0에서는 특히 pyspark에서 많은 지원이 있었다.

      해당 글에서는 이부분 대해서 집중적으로 설명하도록 진행하겠다.

      Spark 3.3.0에서 달라진점

      1. Pandas API on Spark
      2. PySpark의 New Functionalities
      3. PySpark의 Productivity


      1. Pandas API on Spark

      Pandas를 인메모리 분석을 위한 데이터 구조로 제공하도록 지원.



      이미 Apache Spark 3.2.x에서 이미 최초로 지원을 시작했으나 분산연산에 대한 미지원으로 Pandas는 큰 데이터셋을 분석하기에는 다소 까다로왔다. 때문에 Spark 3.3.0에서는 분산시퀀스 인덱싱을 지원함으로써 분산시스템에서 더 나은 확장성을 가지게 되었다


      Python Pandas 에서는 index로 구성되어있으나 Spark 3.2.x의 Pandas에서는 index가 없는 ExistingRDD였다

      이를 Spark 3.3.0에서는 개선하여 Range형태로 distributed-index를 구현하였다.



      아래는 Pyspark Pandas를 import하는 코드이며 index type이 3.3.0에서는 'distributed-sequence'로 변화한 것을 보여준다.

      import pandas as ps #Pandas used
      import pyspark.pandas as ps #Pandas in pyspark used
      ps.options.compute.default_index_type
      'distributed-sequence'

      기존 spark 3.2.x 와 비교했을때 3.3.0의 속도 향상





      초기 Spark 3.2.x에서 지원해주던 pandas API들보다 더 확장된 기능을 지원하도록 추가되었다.

      (Pandas 1.3에서 제공하는 dataframe에서의 기능구현이 이루어짐)

      Supported pandas API


      • ps.merge_asof (SPARK-36813)

      • DataFrame.combine_first (SPARK-36399)

      • DataFrame.cov (SPARK-36396)

      • TimedeltaIndex (SPARK-37525)

      • MultiIndex.dtypes (SPARK-36930)

      • ps.timedelta_range (SPARK-37673)

      • ps.to_timedelta (SPARK-37701)

      • Timedelta Series (SPARK-37525)

      아래 글에서는 각 기능별을 샘플코드로 설명

      ps.merge_asof

      Pandas에서지원하고 있던 merge_asof join을 구현

      from pyspark import pandas as ps
      df1 = ps.DataFrame(
       {"A": ['A1', 'A2', 'A3'],
       "B": [1, 5, 10]})
      df2 = ps.DataFrame(
       {"B": [1, 2, 4, 6, 8],
       "C": ['C1', 'C2', 'C4', 'C6', 'C8']})
      
      result_df = ps.merge_asof(df1, df2, on='B')
       {"A": ['A1', 'A2', 'A3'],
       "B": [1, 5, 10],
       "C": ["C1", "C4", "C8"]}

      DataFrame.combine_first

      Dataframe에서 Null value를 유지로 업데이트 하는 방법은 일반적인 사용사례이므로 지원하도록 업데이트

      from pyspark import pandas as ps
      
      ps.set_option("compute.ops_on_diff_frames", True)
      df1 = ps.DataFrame({'A': [None, 0], 'B': [None, 4]})
      df2 = ps.DataFrame({'A': [1, 1], 'B': [3, 3]})
      df1.combine_first(df2).sort_index()
           A    B
      0  1.0  3.0
      1  0.0  4.0
      
      # Null values still persist if the location of that null value does not exist in other
      
      df1 = ps.DataFrame({'A': [None, 0], 'B': [4, None]})
      df2 = ps.DataFrame({'B': [3, 3], 'C': [1, 1]}, index=[1, 2])
      df1.combine_first(df2).sort_index()
           A    B    C
      0  NaN  4.0  NaN
      1  0.0  3.0  1.0
      2  NaN  3.0  1.0
      ps.reset_option("compute.ops_on_diff_frames")

      DataFrame.cov

      DataFrame에서 Covariance 연산을 지원

      Covariance Matrix로 반환하여 아래와 같이 샘플코드로 구현이 가능해짐

      from pyspark import pandas as ps
      
      psdf = ps.DataFrame([(1, 2), (0, 3), (2, 0), (1, 1)],
                          columns=['dogs', 'cats'])
      psdf.cov()
             dogs      cats
      dogs  0.666667 -1.000000
      cats -1.000000  1.666667
      
      pdf = pd.DataFrame(
              {
                  "a": [1, np.nan, 3, 4],
                  "b": [True, False, False, True],
                  "c": [True, True, False, True],
              }
          )
      psdf = ps.from_pandas(pdf)
      psdf.cov()
                a         b         c
      a  2.333333 -0.166667 -0.166667
      b -0.166667  0.333333  0.166667
      c -0.166667  0.166667  0.250000

      TimedeltaIndex

      pandas에서 timestamp를 제공하는데 있어서 numpy.datetime64 데이터 타입을 기반으로 하는 TimedeltaIndex를 지원한다. 이는 인터벌시간에 대한 정확도나 다양한 기능들을 제공한다.

      from datetime import timedelta
      from pyspark import pandas as ps
      import pandas as pd
      
      ps.from_pandas(
          pd.Series([timedelta(minutes=1)],
          index=pd.TimedeltaIndex([timedelta(days=1)])))
      
      -> 1 days 0 days 00:01:00
      dtype: timedelta64[ns]

      MultiIndex.dtypes

      pandas에서 지원하는 dataType Object로 지원하는 dtype으로 사용하도록 지원

      idx = pd.MultiIndex.from_arrays([[0, 1, 2, 3, 4, 5, 6, 7, 8], [1, 2, 3, 4, 5, 6, 7, 8, 9]], names=("zero", "one"))
      pdf = pd.DataFrame(
          {"a": [1, 2, 3, 4, 5, 6, 7, 8, 9], "b": [4, 5, 6, 3, 2, 1, 0, 0, 0]},
              index=idx,
          )
      psdf = ps.from_pandas(pdf)
      
      ps.DataFrame[psdf.index.dtypes, psdf.dtypes]
      -> typing.Tuple[pyspark.pandas.typedef.typehints.IndexNameType,
          pyspark.pandas.typedef.typehints.IndexNameType,
          pyspark.pandas.typedef.typehints.NameType,
          pyspark.pandas.typedef.typehints.NameType]

      ps.timedelta_range

      pandas의 timedelta_range를 지원, timedelta64단위로 period만큼 TimedeltaIndex를 반환

      from pyspark import pandas as ps
      
      ps.timedelta_range(start="1 day", end="3 days")
      -> TimedeltaIndex(['1 days', '2 days', '3 days'], dtype='timedelta64[ns]', freq=None)
      
      ps.timedelta_range(start="1 day", periods=3)
      -> TimedeltaIndex(['1 days', '2 days', '3 days'], dtype='timedelta64[ns]', freq=None)
      
      ps.timedelta_range(start='1 day', end='2 days', freq='6H')
      -> TimedeltaIndex(['1 days 00:00:00', '1 days 06:00:00', '1 days 12:00:00',
                      '1 days 18:00:00', '2 days 00:00:00'],
                     dtype='timedelta64[ns]', freq=None)

      ps.to_timedelta

      입력받은 str형태의 time value를 timedeltaindex로 변환해주는 함수

      from pyspark import pandas as ps
      
      ps.to_timedelta('1 days 06:05:01.00003')
      -> Timedelta('1 days 06:05:01.000030')
      
      ps.to_timedelta(['1 days 06:05:01.00003', '15.5us', 'nan'])
      -> TimedeltaIndex(['1 days 06:05:01.000030', '0 days 00:00:00.000015500', NaT], dtype='timedelta64[ns]', freq=None)
      
      ps.to_timedelta(pd.Series([1, 2]), unit="d")
      0   1 days
      1   2 days
      dtype: timedelta64[ns]

      2. PySpark의 New Functionalities

      datetime.timedelta support

      datetime.timedelta를 지원

      import datetime
      df = spark.createDataFrame([{'col': datetime.timedelta(minutes=10)}])
      row = df.select(df.col - datetime.timedelta(minutes=9, seconds=12)).first()
      row[0]
      
      -> datetime.timedelta(seconds=48)

      PyArrow batch interface

      PyArrow를 사용에 batch interface를 지원, mapInArrow함수를 통하여 PyArrow 레코드 배치를 하도록 함.

      import pyarrow as pa
      
      df = spark.createDataFrame(
          [(1, "foo"), (2, None), (3, "bar"), (4, "bar")], "a int, b string")
      
      def func(iterator):
          for batch in iterator:
              # `batch` is pyarrow.RecordBatch.
              yield batch
      
      df.mapInArrow(func, df.schema).collect()

      Python standard string formatter in SQL

      Python의 표준 문자열 포맷을 지원하도록 개선

      spark.sql에서 SQL문법을 python포맷으로 지원

      mydf = spark.range(10)
      spark.sql("SELECT {tbl.id}, {tbl[id]} FROM {tbl}", tbl=mydf)
      
      mydf = spark.range(10)
      spark.sql("SELECT * FROM {tbl}", tbl=mydf)
      
      mydf = spark.range(10)
      spark.sql("SELECT {c} FROM {tbl}", c=col("id"), tbl=mydf)
      
      mydf = spark.range(10)
      spark.sql(
          "SELECT {col} FROM {mydf} WHERE id IN {x}",
          col=mydf.id, mydf=mydf, x=tuple(range(4)))

      3. Productivity

      Python/Pandas UDF Profiler

      udf가 적용된 DataFrame의 show_profiles을 제공하여 udf의 코드의 연산계획을 보여준다.

      from pyspark.sql.functions import udf
      import time
      _ = spark.range(10).select(
       udf(lambda x: time.sleep(1))("id")
      ).collect()
      sc.show_profiles()
      
      ============================================================
      Profile of UDF<id=12>
      ============================================================
       30 function calls in 10.009 seconds
       Ordered by: internal time, cumulative time
       ncalls tottime percall cumtime percall filename:lineno(function)
       10 10.009 1.001 10.009 1.001 {built-in method time.sleep}
       10 0.000  0.000 10.009 1.001 <command-1635211799533119>:2(<lambda>)
       10 0.000  0.000  0.000 0.000 {method 'disable' of '_lsprof.Profiler' objects}

      마무리

      spark 3.3.0에서는 pandas의 많은 지원과 업데이트가 있음을 알 수 있었다.

      앞으로 spark next version에서는 아래와 같은 내용을 업데이트하도록 포커싱된다고 한다.



      댓글 0

      DEVOTEE를 활성화 시키면
      지금 작성한 댓글에 AI가 댓글을 달아줍니다.

      jungbbong 님의 최신 블로그

      더보기
      동영상 기고하기