Ошибки Java PySpark EOF и CRCPython

Программы на Python
Anonymous
Ошибки Java PySpark EOF и CRC

Сообщение Anonymous »

Я использую PySpark для обработки данных. Я пробовал это как в Windows 11, так и в WSL2 с Python версии 3.10.12, Java версии 21.0.3 и 17, winutils и Hadoop.dll для Hadoop 3.3.6 и Spark версии 3.5.1. Когда я запускаю свой код, я обычно сталкиваюсь либо с исключением EOF, либо с исключением CRC Java. Есть также несколько других исключений, которые встречаются реже. Возникающая ошибка меняется в зависимости от того, когда я ее запускаю, даже для одного и того же кода. У меня были одни и те же ошибки в другом коде, и мне удалось заставить его работать до конца, только запуская каждую часть в блокнотах Jupyter до тех пор, пока ошибка не исчезла. Вот код:

Код: Выделить всё

from pyspark.sql import SparkSession
from pyspark.sql.functions import udf, explode, col, array_contains, collect_list
from pyspark.sql.types import StructType, StructField, StringType, MapType, IntegerType, ArrayType, BooleanType, FloatType

import nltk
from nltk.stem import WordNetLemmatizer
from spellchecker import SpellChecker
import re
import numpy as np
from numpy.linalg import norm

# Functions from bidir_edge_liwc_mrjob_2 (some are modified, but have the same functionality)
class Parser:
def __init__(self):
# Init code no longer relevant
print('filler')

def count_user_langs(self, timeline):
d = {}
for i in timeline:
try:
d[i] += 1
except KeyError:
d[i] = 1
return d

def _n_most_values(self, d, n_values, proportion_threshold):
if 'und' in d: # 'und' is undefined (language)
del d['und']
if '' in d:
del d['']
sum1 = sum(d.values()) #number of tweets with location information.
if len( d ) > 0:
attrib = []
t=0
for w in sorted(d, key=d.get, reverse=True): #sort keys in descending "value"
if n_values == 1: #location with maximum frequency
attrib.append(w)
break
else:
if t < n_values:
if t==0:
attrib.append(w)
else:
if d[w] / float(sum1) > proportion_threshold:
attrib.append(w)
t+=1
return attrib

else: #If there was no location information
return []

class Vectorizer:
def __init__(self):
nltk.download('stopwords')
nltk.download('wordnet')
nltk.download('omw-1.4')
nltk.download('punkt')

self.lemma = WordNetLemmatizer()
self.spell = SpellChecker()
self.vector_path = 'C:\\Users\\jjhen\\Personal Documents\\research - Copy\\glove_twitter\\glove.twitter.27B.200d.txt'

def remove_unwanted(self, document):
# remove email addresses
document = re.sub(r"([a-zA-Z0-9_.+-]+@[a-zA-Z0-9-]+\.[a-zA-Z0-9-.]+)", '', document)
# remove user mentions
document = re.sub("@[A-Za-z0-9_]+","", document)
# remove URLS
document = re.sub(r'http\S+', '', document)
# remove hashtags
document = re.sub("#[A-Za-z0-9_]+","", document)
# remove unwanted characters
document = re.sub("[^0-9A-Za-z ]", ""  , document)
# remove multiple spaces
document = document.replace(' +',' ')
return document.strip()

def remove_words(self, tokens):
stopwords = nltk.corpus.stopwords.words('english')
# remove puntcuation from stopwords
stopwords = [self.remove_unwanted(word) for word in stopwords]
cleaned_tokens = [token for token in tokens if token not in stopwords]
return cleaned_tokens

def lemmatize(self, tokens):
lemmatized_tokens = [self.lemma.lemmatize(token, pos = 'v') for token in tokens]
return lemmatized_tokens

def preprocess(self, document):
document = document.lower()
document = re.sub(r'\n', '', document)
document = self.remove_unwanted(document)
tokens = document.split()
tokens = self.remove_words(tokens)
# tokens = correct(tokens)
tokens = self.lemmatize(tokens)
return tokens

# Returns vector as an array for Spark
def load_glove_vectors(self, path=None):
if path is None:
path = self.vector_path
embeddings = {}
with open(path, 'r', encoding='utf-8') as f:
for line in f:
parts = line.strip().split()
word = parts[0]
vector = np.array(parts[1:], dtype=float)
embeddings[word] = vector
return embeddings

# TODO handle case where no words
def get_average_vector(self, document, embeddings):
# tokens = self.preprocess(document)
# print(tokens)
# vectors = [embeddings[token] for token in tokens if token in embeddings]
# return np.array(np.mean(vectors, axis=0)).tolist()
return [1.0,2.0,3.0]

def cosine_similarity(self, vec_1, vec_2):
return np.dot(vec_1, vec_2) / (norm(vec_1) * norm(vec_2))

vectorizer = Vectorizer()
embeddings = vectorizer.load_glove_vectors()
print(vectorizer.get_average_vector('', embeddings))
spark = SparkSession.builder.appName('spark_parse').config('spark.network.timeout', '800s').config('spark.network.timeout', '800s').config('spark.executor.heartbeatInterval', '120s').config("spark.executor.memory", "30g").config("spark.driver.memory", "30g").getOrCreate()
spark.conf.set('spark.sql.debug.maxToStringFields', 1000)
spark.sparkContext.setLogLevel("WARN")

@udf(returnType = StringType())
def get_first_udf(list):
return list[0]

# Change these paths to match your directory location
tweets_path = 'C:\\Users\\jjhen\\Personal Documents\\research - Copy\\NZ'
triads_path = 'C:\\Users\\jjhen\\Personal Documents\\research - Copy\\Triads_NZ'

# Read and format tweets, filter by language, ensure dominant language, and exclude retweets
tweets_df = spark.read.option('compression', 'bzip2').option('header', 'true').json(tweets_path)
# tweets_df = tweets_df.filter(has_dom_lang_udf(col('tweets.lang')))
tweets_df = tweets_df.select(
col('user.id_str').alias('user_id'),
explode(col('tweets')).alias('tweets'),
)
tweets_df = tweets_df.withColumn('mentioned_users', col('tweets.entities.user_mentions.id_str'))
tweets_df = tweets_df.withColumn('user_id', get_first_udf(col('user_id')))
tweets_df = tweets_df.filter((~col('tweets.text').contains('RT')) &
((col('tweets.lang') == 'en') |
(col('tweets.lang') == 'und')))

# Read and format
triads_df = spark.read.option('compression', 'bzip2').option('delimiter', '\t').option('header', 'true').csv(triads_path)
triads = triads_df.select(triads_df.columns[:3])

# Match tweets with dyads (user is dyad's node_1)
joined_df1 = tweets_df.join(
triads_df,
(tweets_df.user_id == triads_df.node_1) &
((array_contains(tweets_df.mentioned_users, triads_df.node_2)) |
(col('tweets.in_reply_to_user_id_str') == triads_df.node_2)) &
((array_contains(tweets_df.mentioned_users, triads_df.node_3)) |
(col('tweets.in_reply_to_user_id_str') == triads_df.node_3)),
'inner'
)

# Match tweets with dyads (user is dyad's node_2)
joined_df2 = tweets_df.join(
triads_df,
(tweets_df.user_id == triads_df.node_2) &
((array_contains(tweets_df.mentioned_users, triads_df.node_1)) |
(col('tweets.in_reply_to_user_id_str') == triads_df.node_1)) &
((array_contains(tweets_df.mentioned_users,  triads_df.node_3)) |
(col('tweets.in_reply_to_user_id_str') == triads_df.node_3)),
'inner'
)

# Match tweets with dyads (user is dyad's node_3)
joined_df3 = tweets_df.join(
triads_df,
(tweets_df.user_id == triads_df.node_3) &
((array_contains(tweets_df.mentioned_users, triads_df.node_1)) |
(col('tweets.in_reply_to_user_id_str') == triads_df.node_1)) &
((array_contains(tweets_df.mentioned_users, triads_df.node_2)) |
(col('tweets.in_reply_to_user_id_str') == triads_df.node_2)),
'inner'
)

union_df = joined_df1.union(joined_df2).union(joined_df3)

# Group by triad and collect LIWC info
grouped_df = union_df.groupBy('node_1', 'node_2', 'node_3').agg(
collect_list('tweets.text').alias('tweets')
)

def append_lists(lists):
return ' '.join(lists)

# Aggregate LIWC data for dyad by summing values for each LIWC category
append_udf = udf(append_lists, StringType())
grouped_df = grouped_df.withColumn('tweets', append_udf(col('tweets')))

vectorize_udf = udf(lambda tweet: vectorizer.get_average_vector(tweet, embeddings), ArrayType(FloatType()))
vector_udf = grouped_df.select(col('node_1'), col('node_2'), col('node_3'), vectorize_udf(col('tweets')).alias('vector'))
vector_udf.write.option('header', 'true').csv('C:\\Users\\jjhen\\Personal Documents\\research - Copy\\vectors')
Ниже приведены трассировки стека для двух наиболее распространенных ошибок:

Код: Выделить всё

java.io.IOException:  crc error
at org.apache.hadoop.io.compress.bzip2.CBZip2InputStream.reportCRCError(CBZip2InputStream.java:258)
at org.apache.hadoop.io.compress.bzip2.CBZip2InputStream.endBlock(CBZip2InputStream.java:569)
at org.apache.hadoop.io.compress.bzip2.CBZip2InputStream.setupNoRandPartA(CBZip2InputStream.java:1104)
at org.apache.hadoop.io.compress.bzip2.CBZip2InputStream.setupNoRandPartB(CBZip2InputStream.java:1158)
at org.apache.hadoop.io.compress.bzip2.CBZip2InputStream.read0(CBZip2InputStream.java:464)
at org.apache.hadoop.io.compress.bzip2.CBZip2InputStream.read(CBZip2InputStream.java:419)
at org.apache.hadoop.io.compress.BZip2Codec$BZip2CompressionInputStream.read(BZip2Codec.java:490)
at java.base/java.io.InputStream.read(InputStream.java:218)
at org.apache.hadoop.mapreduce.lib.input.CompressedSplitLineReader.fillBuffer(CompressedSplitLineReader.java:130)
at org.apache.hadoop.util.LineReader.readDefaultLine(LineReader.java:227)
at org.apache.hadoop.util.LineReader.readLine(LineReader.java:185)
at org.apache.hadoop.mapreduce.lib.input.CompressedSplitLineReader.readLine(CompressedSplitLineReader.java:159)
at org.apache.hadoop.mapreduce.lib.input.LineRecordReader.nextKeyValue(LineRecordReader.java:200)
at org.apache.spark.sql.execution.datasources.RecordReaderIterator.hasNext(RecordReaderIterator.scala:39)
at org.apache.spark.sql.execution.datasources.HadoopFileLinesReader.hasNext(HadoopFileLinesReader.scala:67)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.hasNext(FileScanRDD.scala:129)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:491)
at scala.collection.Iterator.foreach(Iterator.scala:943)
at scala.collection.Iterator.foreach$(Iterator.scala:943)
at scala.collection.AbstractIterator.foreach(Iterator.scala:1431)
at scala.collection.TraversableOnce.reduceLeft(TraversableOnce.scala:237)
at scala.collection.TraversableOnce.reduceLeft$(TraversableOnce.scala:220)
at scala.collection.AbstractIterator.reduceLeft(Iterator.scala:1431)
at scala.collection.TraversableOnce.reduceLeftOption(TraversableOnce.scala:249)
at scala.collection.TraversableOnce.reduceLeftOption$(TraversableOnce.scala:248)
at scala.collection.AbstractIterator.reduceLeftOption(Iterator.scala:1431)
at scala.collection.TraversableOnce.reduceOption(TraversableOnce.scala:256)
at scala.collection.TraversableOnce.reduceOption$(TraversableOnce.scala:256)
at scala.collection.AbstractIterator.reduceOption(Iterator.scala:1431)
at org.apache.spark.sql.catalyst.json.JsonInferSchema.$anonfun$infer$1(JsonInferSchema.scala:107)
at org.apache.spark.rdd.RDD.$anonfun$mapPartitions$2(RDD.scala:858)
at org.apache.spark.rdd.RDD.$anonfun$mapPartitions$2$adapted(RDD.scala:858)
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93)
at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:166)
at org.apache.spark.scheduler.Task.run(Task.scala:141)
at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620)
at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64)
at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61)
at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
at java.base/java.lang.Thread.run(Thread.java:842)
Вторая ошибка:

Код: Выделить всё

ERROR Executor: Exception in task 0.0 in stage 5.0 (TID 10483)
org.apache.spark.SparkException:  Python worker exited unexpectedly (crashed)
at org.apache.spark.api.python.BasePythonRunner$ReaderIterator$$anonfun$1.applyOrElse(PythonRunner.scala:612)
at org.apache.spark.api.python.BasePythonRunner$ReaderIterator$$anonfun$1.applyOrElse(PythonRunner.scala:594)
at scala.runtime.AbstractPartialFunction.apply(AbstractPartialFunction.scala:38)
at org.apache.spark.sql.execution.python.BasePythonUDFRunner$$anon$1.read(PythonUDFRunner.scala:99)
at org.apache.spark.sql.execution.python.BasePythonUDFRunner$$anon$1.read(PythonUDFRunner.scala:75)
at org.apache.spark.api.python.BasePythonRunner$ReaderIterator.hasNext(PythonRunner.scala:525)
at org.apache.spark.InterruptibleIterator.hasNext(InterruptibleIterator.scala:37)
at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:491)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage2.processNext(Unknown Source)
at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
at org.apache.spark.sql.execution.WholeStageCodegenEvaluatorFactory$WholeStageCodegenPartitionEvaluator$$anon$1.hasNext(WholeStageCodegenEvaluatorFactory.scala:43)
at org.apache.spark.sql.execution.SparkPlan.$anonfun$getByteArrayRdd$1(SparkPlan.scala:388)
at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:893)
at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:893)
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93)
at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:166)
at org.apache.spark.scheduler.Task.run(Task.scala:141)
at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620)
at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64)
at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61)
at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
at java.base/java.lang.Thread.run(Thread.java:842)
Caused by: java.io.EOFException
at java.base/java.io.DataInputStream.readInt(DataInputStream.java:398)
at org.apache.spark.sql.execution.python.BasePythonUDFRunner$$anon$1.read(PythonUDFRunner.scala:83)
...  26 more
24/07/02 19:41:46 WARN TaskSetManager: Lost task 0.0 in stage 5.0 (TID 10483) (192.168.0.34 executor driver): org.apache.spark.SparkException: Python worker exited unexpectedly (crashed)
at org.apache.spark.api.python.BasePythonRunner$ReaderIterator$$anonfun$1.applyOrElse(PythonRunner.scala:612)
at org.apache.spark.api.python.BasePythonRunner$ReaderIterator$$anonfun$1.applyOrElse(PythonRunner.scala:594)
at scala.runtime.AbstractPartialFunction.apply(AbstractPartialFunction.scala:38)
at org.apache.spark.sql.execution.python.BasePythonUDFRunner$$anon$1.read(PythonUDFRunner.scala:99)
at org.apache.spark.sql.execution.python.BasePythonUDFRunner$$anon$1.read(PythonUDFRunner.scala:75)
at org.apache.spark.api.python.BasePythonRunner$ReaderIterator.hasNext(PythonRunner.scala:525)
at org.apache.spark.InterruptibleIterator.hasNext(InterruptibleIterator.scala:37)
at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:491)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage2.processNext(Unknown Source)
at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
at org.apache.spark.sql.execution.WholeStageCodegenEvaluatorFactory$WholeStageCodegenPartitionEvaluator$$anon$1.hasNext(WholeStageCodegenEvaluatorFactory.scala:43)
at org.apache.spark.sql.execution.SparkPlan.$anonfun$getByteArrayRdd$1(SparkPlan.scala:388)
at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:893)
at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:893)
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93)
at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:166)
at org.apache.spark.scheduler.Task.run(Task.scala:141)
at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620)
at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64)
at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61)
at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
at java.base/java.lang.Thread.run(Thread.java:842)
Caused by: java.io.EOFException
at java.base/java.io.DataInputStream.readInt(DataInputStream.java:398)
at org.apache.spark.sql.execution.python.BasePythonUDFRunner$$anon$1.read(PythonUDFRunner.scala:83)
...  26 more
...
Третья ошибка:

Код: Выделить всё

ERROR Executor: Exception in task 187.0 in stage 1.0 (TID 10187)
java.lang.ArrayIndexOutOfBoundsException: Index 134886018 out of bounds for length 900000
at org.apache.hadoop.io.compress.bzip2.CBZip2InputStream.setupNoRandPartA(CBZip2InputStream.java:1095)
at org.apache.hadoop.io.compress.bzip2.CBZip2InputStream.setupNoRandPartB(CBZip2InputStream.java:1158)
at org.apache.hadoop.io.compress.bzip2.CBZip2InputStream.read0(CBZip2InputStream.java:464)
at org.apache.hadoop.io.compress.bzip2.CBZip2InputStream.read(CBZip2InputStream.java:419)
at org.apache.hadoop.io.compress.BZip2Codec$BZip2CompressionInputStream.read(BZip2Codec.java:490)
at java.base/java.io.InputStream.read(InputStream.java:218)
at org.apache.hadoop.mapreduce.lib.input.CompressedSplitLineReader.fillBuffer(CompressedSplitLineReader.java:130)
at org.apache.hadoop.util.LineReader.readDefaultLine(LineReader.java:227)
at org.apache.hadoop.util.LineReader.readLine(LineReader.java:185)
at org.apache.hadoop.mapreduce.lib.input.CompressedSplitLineReader.readLine(CompressedSplitLineReader.java:159)
at org.apache.hadoop.mapreduce.lib.input.LineRecordReader.nextKeyValue(LineRecordReader.java:200)
at org.apache.spark.sql.execution.datasources.RecordReaderIterator.hasNext(RecordReaderIterator.scala:39)
at ...
В других сообщениях на форумах в Интернете указывается, что это может быть причиной прерываний из-за нехватки ресурсов. Я установил файл spark-defaults.conf, чтобы предоставить 30 г памяти как драйверу, так и исполнителю, и я также настроил это в коде. Я также отслеживал использование памяти во время работы (просто с помощью диспетчера задач), и, похоже, он не приближается к максимальному объему памяти моего компьютера в 32 ГБ (при работе в Windows он приближается к WSL).Я также попытался убедиться, что ни один из сжатых файлов .json.bz2 в папкеtwitters_path не был поврежден. Я использовал:

Код: Выделить всё

bzip2 -tv *.bz2
чтобы убедиться, что все файлы в порядке. Я также написал скрипт Python, чтобы проверить это.
Я пробовал работать на WSL и Windows с несколькими различными версиями Java. Я также использовал несколько разных источников winutils.exe и Hadoop.dll, чтобы убедиться, что это не ошибка конфигурации.
Буду признателен за любую помощь.

Подробнее здесь: https://stackoverflow.com/questions/786 ... ava-errors

Вернуться в «Python»