SparkConnectGrpcException при работе с DataframePython

Программы на Python
Anonymous
SparkConnectGrpcException при работе с Dataframe

Сообщение Anonymous »


Я использую Spark Connect для подключения к удаленному кластеру Spark и выполнения этого блока кода. Кажется, я что-то упускаю. Вот блокнот IPython. Я использую PySpark, и для запуска удаленного сервера я использую эту команду на стороне сервера:

/opt/bitnami/spark/sbin/start-connect-server.sh --packages org.apache.spark:spark-connect_2.12:3.5.0 --master spark://spark-master-0.spark-headless.spark-operator.svc.cluster.local:7077 Блокнот:

импортировать панд как pd импортировать numpy как np от даты и времени импорта, timedelta, datetime время импорта импортировать pyspark из pyspark.sql импортировать SparkSession, SQLContext из pyspark.context импортировать SparkContext из импорта pyspark.sql.functions * из импорта pyspark.sql.types * из даты и времени импорта даты и времени, даты из строки импорта pyspark.sql из pyspark.sql импортировать SparkSession из импорта pyspark.sql.functions * spark = SparkSession.builder.appName("Мое приложение").remote("sc://HOST:PORT").getOrCreate() df = spark.createDataFrame([ Row(a=1, b=2., c='string1', d=date(2000, 1, 1), e=datetime(2000, 1, 1, 12, 0)), Row(a=2, b=3., c='string2', d=date(2000, 2, 1), e=datetime(2000, 1, 2, 12, 0)), Row(a=4, b=5., c='string3', d=date(2000, 3, 1), e=datetime(2000, 1, 3, 12, 0)) ]) df.collect() Вывод: [Row(a=1, b=2.0, c='string1', d=datetime.date(2000, 1, 1), e=datetime.datetime(2000, 1, 1, 12, 0)), Row(a=2, b=3.0, c='string2', d=datetime.date(2000, 2, 1), e=datetime.datetime(2000, 1, 2, 12, 0)), Row(a=4, b=5.0, c='string3', d=datetime.date(2000, 3, 1), e=datetime.datetime(2000, 1, 3, 12, 0))] df = spark.read.format('csv').option('header','true').load('breast-cancer.csv') df2 = df.select(df.columns[0:4]) df2.show(n=5) Выход:

+--------+---------+-----------+------------ + | id|диагноз|среднее_радиуса|среднее_текстуры| +--------+---------+-----------+------------+ | 842302| М | 17,99 | 10.38| | 842517| М | 20.57| 17,77 | |84300903| М | 19,69 | 21.25| |84348301| М | 11,42 | 20.38| |84358402| М | 20.29| 14,34 | +--------+---------+-----------+------------+ показаны только 5 верхних строк

Но когда я это сделаю:

df = spark.read.format('csv').option('header','true').load('breast-cancer.csv') df.collect() или,

df = spark.read.format('csv').option('header','true').load('breast-cancer.csv') df.limit(5).toPandas() Выход:

-------------------------------------------- ------------------------------- SparkConnectGrpcException Traceback (самый последний вызов — последний) Ячейка In[7], строка 2 1 df = spark.read.format('csv').option('header','true').load('breast-cancer.csv') ----> 2 df.collect() Файл ~/Documents/workstation/jupyter-lab/venv/lib/python3.10/site-packages/pyspark/sql/connect/dataframe.py:1645 в DataFrame.collect(self) 1643 поднять исключение («Невозможно собрать данные в пустом сеансе»). 1644 запрос = self._plan.to_proto(self._session.client) -> Таблица 1645, схема = self._session.client.to_table(запрос) 1647 схема = схема или from_arrow_schema(table.schema, предпочитает_timestamp_ntz=True) 1649 схема утверждения не None, а isinstance(schema, StructType) Файл ~/Documents/workstation/jupyter-lab/venv/lib/python3.10/site-packages/pyspark/sql/connect/client/core.py:858 в SparkConnectClient.to_table(self, plan) 856 req = self._execute_plan_request_with_metadata() 857 req.plan.CopyFrom(план) --> Таблица 858, схема, _, _, _ = self._execute_and_fetch(req) 859 таблица утверждений не имеет значения None Таблица возврата 860, схема Файл ~/Documents/workstation/jupyter-lab/venv/lib/python3.10/site-packages/pyspark/sql/connect/client/core.py:1282 в SparkConnectClient._execute_and_fetch(self, req, self_destruct) Схема 1279: Необязательно [StructType] = Нет 1280 свойств: Dict[str, Any] = {} -> 1282 для ответа в self._execute_and_fetch_as_iterator(req): 1283, если isinstance(ответ, StructType): 1284 схема = ответ Файл ~/Documents/workstation/jupyter-lab/venv/lib/python3.10/site-packages/pyspark/sql/connect/client/core.py:1263 в SparkConnectClient._execute_and_fetch_as_iterator(self, req) 1261 выход из handle_response(b) 1262, кроме исключения как ошибки: -> 1263 self._handle_error(ошибка) Файл ~/Documents/workstation/jupyter-lab/venv/lib/python3.10/site-packages/pyspark/sql/connect/client/core.py:1502 в SparkConnectClient._handle_error(self, error) 1489 """ 1490 Обработка ошибок, возникающих во время вызовов RPC. 1491 (...) 1499 Вызывает соответствующее внутреннее исключение Python. 1500 """ 1501, если isinstance(ошибка, grpc.RpcError): -> 1502 self._handle_rpc_error(ошибка) 1503 elif isinstance (ошибка, ValueError): 1504, если «Невозможно вызвать RPC» в строке (ошибка) и «закрыто» в строке (ошибка): Файл ~/Documents/workstation/jupyter-lab/venv/lib/python3.10/site-packages/pyspark/sql/connect/client/core.py:1538 в SparkConnectClient._handle_rpc_error(self, rpc_error) 1536 информация = error_details_pb2.ErrorInfo() 1537 г.Распаковать(информация) -> 1538 поднять Convert_Exception(info, status.message) с None 1540 поднять SparkConnectGrpcException(status.message) с None 1541 еще: SparkConnectGrpcException: (org.apache.spark.SparkException) Задание прервано из-за сбоя этапа: задача 0 на этапе 62.0 завершилась неудачно 4 раза, последний сбой: потеряна задача 0,3 на этапе 62.0 (TID 106) (172.31.30.154 исполнитель 0): java .lang.ClassCastException: невозможно назначить экземпляр java.lang.invoke.SerializedLambda полю org.apache.spark.rdd.MapPartitionsRDD.f типа scala.Function3 в экземпляре org.apache.spark.rdd.MapPartitionsRDD в java.base/java.io.ObjectStreamClass$FieldReflector.setObjFieldValues(ObjectStreamClass.java:2096) в java.base/java.io.ObjectStreamClass$FieldReflector.checkObjectFieldValueTypes(ObjectStreamClass.java:2060) в java.base/java.io.ObjectStreamClass.checkObjFieldValueTypes(ObjectStreamClass.java:1347) в java.base/java.io.ObjectInputStream$FieldValues.defaultCheckFieldValues(ObjectInputStream.java:2679) в java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2486) в java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257) в java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733) в java.base/java.io.ObjectInputStream$FieldValues.(ObjectInputStream.java:2606) в java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457) в java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257) в java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733) в java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:509) в java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:467) в org.apache.spark.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:87) в org.apache.spark.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:129) в org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:86) в org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:161) в org.apache.spark.scheduler.Task.run(Task.scala:141) в org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620) в org.apache.spark.util.SparkE... Как решить эту проблему?

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