Estaba usando pyspark en AWS EMR (4 r5.xlarge como 4 trabajadores, cada uno tiene un ejecutor y 4 núcleos), y obtuve AttributeError: Can't get attribute 'new_block' on <module 'pandas.core.internals.blocks' . A continuación se muestra un fragmento del código que arrojó este error:
search = SearchEngine(db_file_dir = "/tmp/db") conn = sqlite3.connect("/tmp/db/simple_db.sqlite") pdf_ = pd.read_sql_query('''select zipcode, lat, lng, bounds_west, bounds_east, bounds_north, bounds_south from simple_zipcode''',conn) brd_pdf = spark.sparkContext.broadcast(pdf_) conn.close() @udf('string') def get_zip_b(lat, lng): pdf = brd_pdf.value out = pdf[(np.array(pdf["bounds_north"]) >= lat) & (np.array(pdf["bounds_south"]) <= lat) & (np.array(pdf['bounds_west']) <= lng) & (np.array(pdf['bounds_east']) >= lng) ] if len(out): min_index = np.argmin( (np.array(out["lat"]) - lat)**2 + (np.array(out["lng"]) - lng)**2) zip_ = str(out["zipcode"].iloc[min_index]) else: zip_ = 'bad' return zip_ df = df.withColumn('zipcode', get_zip_b(col("latitude"),col("longitude"))) A continuación se muestra el rastreo, donde la línea 102, en get_zip_b se refiere a pdf = brd_pdf.value :
21/08/02 06:18:19 WARN TaskSetManager: Lost task 12.0 in stage 7.0 (TID 1814, ip-10-22-17-94.pclc0.merkle.local, executor 6): org.apache.spark.api.python.PythonException: Traceback (most recent call last): File "/mnt/yarn/usercache/hadoop/appcache/application_1627867699893_0001/container_1627867699893_0001_01_000009/pyspark.zip/pyspark/worker.py", line 605, in main process() File "/mnt/yarn/usercache/hadoop/appcache/application_1627867699893_0001/container_1627867699893_0001_01_000009/pyspark.zip/pyspark/worker.py", line 597, in process serializer.dump_stream(out_iter, outfile) File "/mnt/yarn/usercache/hadoop/appcache/application_1627867699893_0001/container_1627867699893_0001_01_000009/pyspark.zip/pyspark/serializers.py", line 223, in dump_stream self.serializer.dump_stream(self._batched(iterator), stream) File "/mnt/yarn/usercache/hadoop/appcache/application_1627867699893_0001/container_1627867699893_0001_01_000009/pyspark.zip/pyspark/serializers.py", line 141, in dump_stream for obj in iterator: File "/mnt/yarn/usercache/hadoop/appcache/application_1627867699893_0001/container_1627867699893_0001_01_000009/pyspark.zip/pyspark/serializers.py", line 212, in _batched for item in iterator: File "/mnt/yarn/usercache/hadoop/appcache/application_1627867699893_0001/container_1627867699893_0001_01_000009/pyspark.zip/pyspark/worker.py", line 450, in mapper result = tuple(f(*[a[o] for o in arg_offsets]) for (arg_offsets, f) in udfs) File "/mnt/yarn/usercache/hadoop/appcache/application_1627867699893_0001/container_1627867699893_0001_01_000009/pyspark.zip/pyspark/worker.py", line 450, in <genexpr> result = tuple(f(*[a[o] for o in arg_offsets]) for (arg_offsets, f) in udfs) File "/mnt/yarn/usercache/hadoop/appcache/application_1627867699893_0001/container_1627867699893_0001_01_000009/pyspark.zip/pyspark/worker.py", line 90, in <lambda> return lambda *a: f(*a) File "/mnt/yarn/usercache/hadoop/appcache/application_1627867699893_0001/container_1627867699893_0001_01_000009/pyspark.zip/pyspark/util.py", line 121, in wrapper return f(*args, **kwargs) File "/mnt/var/lib/hadoop/steps/s-1IBFS0SYWA19Z/Mobile_ID_process_center.py", line 102, in get_zip_b File "/mnt/yarn/usercache/hadoop/appcache/application_1627867699893_0001/container_1627867699893_0001_01_000009/pyspark.zip/pyspark/broadcast.py", line 146, in value self._value = self.load_from_path(self._path) File "/mnt/yarn/usercache/hadoop/appcache/application_1627867699893_0001/container_1627867699893_0001_01_000009/pyspark.zip/pyspark/broadcast.py", line 123, in load_from_path return self.load(f) File "/mnt/yarn/usercache/hadoop/appcache/application_1627867699893_0001/container_1627867699893_0001_01_000009/pyspark.zip/pyspark/broadcast.py", line 129, in load return pickle.load(file) AttributeError: Can't get attribute 'new_block' on <module 'pandas.core.internals.blocks' from '/mnt/miniconda/lib/python3.9/site-packages/pandas/core/internals/blocks.py'>Algunas observaciones y proceso de pensamiento:
1, después de realizar una búsqueda en línea, ¿el AttributeError en pyspark parece ser causado por versiones de pandas que no coinciden entre el controlador y los trabajadores?
2, pero ejecuté el mismo código en dos conjuntos de datos diferentes, uno funcionó sin errores pero el otro no, lo que parece muy extraño e indeterminado, y parece que los errores pueden no ser causados por versiones de pandas que no coinciden. De lo contrario, ninguno de los dos conjuntos de datos tendría éxito.
3, luego ejecuté el mismo código en el conjunto de datos exitoso nuevamente, pero esta vez con diferentes configuraciones de chispa: configurando spark.driver.memory de 2048M a 4192m, y arrojó AttributeError.
4, en conclusión, creo que AttributeError tiene algo que ver con el controlador. Pero no puedo decir cómo están relacionados con el mensaje de error y cómo solucionarlo: AttributeError: no se puede obtener el atributo 'new_block' en <módulo 'pandas.core.internals.blocks'.