EMR을 통해 s3a에 쓸 때 OutOfMemory 오류

Aug 30 2020

다음 PySpark 코드에 대해 OutOfMemory 오류가 발생했습니다. (특정 수의 행이 기록 된 후 실패합니다. s3a를 사용하는 대신 hadoop 파일 시스템에 쓰려고하면이 문제가 발생하지 않으므로 다음으로 좁혔다 고 생각합니다. 문제는 s3a입니다.)-s3a에 쓰는 최종 목표. 매우 큰 테이블에 대해 메모리가 부족하지 않는 최적의 s3a 구성이 있는지 궁금합니다.

df = spark.sql("SELECT * FROM my_big_table")
df.repartition(1).write.option("header", "true").csv("s3a://mycsvlocation/folder/")

내 s3a 구성 (emr 기본값) :

('fs.s3a.attempts.maximum', '10')
('fs.s3a.buffer.dir', '${hadoop.tmp.dir}/s3a')
('fs.s3a.connection.establish.timeout', '5000')
('fs.s3a.connection.maximum', '15')
('fs.s3a.connection.ssl.enabled', 'true')
('fs.s3a.connection.timeout', '50000')
('fs.s3a.fast.buffer.size', '1048576')
('fs.s3a.fast.upload', 'true')
('fs.s3a.impl', 'org.apache.hadoop.fs.s3a.S3AFileSystem')
('fs.s3a.max.total.tasks', '1000')
('fs.s3a.multipart.purge', 'false')
('fs.s3a.multipart.purge.age', '86400')
('fs.s3a.multipart.size', '104857600')
('fs.s3a.multipart.threshold', '2147483647')
('fs.s3a.paging.maximum', '5000')
('fs.s3a.threads.core', '15')
('fs.s3a.threads.keepalivetime', '60')
('fs.s3a.threads.max', '256')
('mapreduce.fileoutputcommitter.algorithm.version', '2')
('spark.authenticate', 'true')
('spark.network.crypto.enabled', 'true')
('spark.network.crypto.saslFallback', 'true')
('spark.speculation', 'false')

스택 추적의 기본 :

Caused by: java.lang.OutOfMemoryError
        at java.io.ByteArrayOutputStream.hugeCapacity(ByteArrayOutputStream.java:123)
        at java.io.ByteArrayOutputStream.grow(ByteArrayOutputStream.java:117)
        at java.io.ByteArrayOutputStream.ensureCapacity(ByteArrayOutputStream.java:93)
        at java.io.ByteArrayOutputStream.write(ByteArrayOutputStream.java:153)
        at org.apache.hadoop.fs.s3a.S3AFastOutputStream.write(S3AFastOutputStream.java:194)
        at org.apache.hadoop.fs.FSDataOutputStream$PositionCache.write(FSDataOutputStream.java:60)
        at java.io.DataOutputStream.write(DataOutputStream.java:107)
        at sun.nio.cs.StreamEncoder.writeBytes(StreamEncoder.java:221)
        at sun.nio.cs.StreamEncoder.implWrite(StreamEncoder.java:282)
        at sun.nio.cs.StreamEncoder.write(StreamEncoder.java:125)
        at java.io.OutputStreamWriter.write(OutputStreamWriter.java:207)
        at com.univocity.parsers.common.input.WriterCharAppender.writeCharsAndReset(WriterCharAppender.java:152)
        at com.univocity.parsers.common.AbstractWriter.writeRow(AbstractWriter.java:808)
        ... 16 more

답변

1 ProgrammingUnicorn Sep 01 2020 at 10:07

여기서 문제는 기본 s3a 업로드가 2GB 또는 2147483647 바이트보다 큰 단일 대용량 파일의 업로드를 지원하지 않는다는 것입니다.

('fs.s3a.multipart.threshold', '2147483647')

내 EMR 버전이 최신 버전보다 오래되었으므로 multipart.threshold 매개 변수는 정수일 뿐이므로 단일 "부분"또는 파일에 대한 제한은 2147483647 바이트입니다. 최신 버전은 int 대신 long을 사용하며 더 큰 단일 파일 크기 제한을 지원할 수 있습니다.

파일을 로컬 hdfs에 쓴 다음 별도의 Java 프로그램을 통해 s3로 이동하는 방법을 사용할 것입니다.