23.06.15
DEVOTEE를 활성화 시키면
지금 작성한 커뮤니티 글에 대해 1개의 댓글을 달아줍니다.
버튼을 누르면 글 수정 시 ChatGPT가 작성한 댓글이 수정됩니다.
| 컨텐츠 유형 | 제목 | 저장일 | 삭제 |
|---|
본인인증 로그인에 실패하였습니다.
회원이 아니시거나 본인인증 등록이
완료되지 않은 사용자입니다.
*해당 글은 spark 3.3.0 릴리즈노트의 1608개의 JIRA에서 개선된 내용을 요약 정리 및 Data+AI Summit 2022에서 발표한 내용들을 합쳐서 만든 글입니다.
최근 spark은 python사용자들이 늘어남에 따라서 python에 대한 지원을 좀 더 확장하도록 프로젝트가 신설(Project Zen)되어서 진행중이다. 해당 프로젝트의 결과로 이번에 릴리즈된 spark 3.3.0에서는 특히 pyspark에서 많은 지원이 있었다.
해당 글에서는 이부분 대해서 집중적으로 설명하도록 진행하겠다.
1. Pandas API on Spark
2. PySpark의 New Functionalities
3. PySpark의 ProductivityPandas를 인메모리 분석을 위한 데이터 구조로 제공하도록 지원.
이미 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에서의 기능구현이 이루어짐)
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)
아래 글에서는 각 기능별을 샘플코드로 설명
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에서 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에서 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.250000pandas에서 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]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]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)입력받은 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]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를 지원, 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의 표준 문자열 포맷을 지원하도록 개선
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)))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에서는 아래와 같은 내용을 업데이트하도록 포커싱된다고 한다.
DEVOTEE를 활성화 시키면
지금 작성한 댓글에 AI가 댓글을 달아줍니다.