Файл метаданных может быть таким
Код: Выделить всё
+--------------------+-------------------+
| dwh_column_name| dwh_data_type|
+--------------------+-------------------+
| product_id | string|
| phone_num| int|
| phone_format_code| string|
| phone_num_type| string|
+--------------------+-------------------+
Код: Выделить всё
+-----------+-----------------+----------+-------------------+
|product_id |phone_format_code| phone_num| phone_num_type|
+-----------+-----------------+----------+-------------------+
| 3241 | (001)706-2887|5829545023| Home Cell Phone|
| 43df | (391)210-4894|988nn09488| Home Fax Number|
| 45fg | (001)202-8944|37aa501564| Business Phone|
с помощью файла метаданных мне нужно проверить типы данных каждого из них значения, если данные соответствуют ожиданиям, мне нужно добавить их в действительный фрейм данных и предположим, что если значения не такие, как ожидалось, мне нужно добавить их в недопустимый фрейм данных.
мне это нужно нужно сделать в Py-spark
мне нужно получить один и тот же тип данных в допустимом фрейме данных
и недопустимый тип данных всей строки в недопустимом фрейме данных
все эти преобразования необходимо выполнить в py-spark
Код: Выделить всё
from pyspark.sql.functions import col
def type_check(valid_range_df, filter_table):
# Извлекаем допустимые типы данных из valid_range_dfvalid_data_types = dict(valid_range_df.dtypes)
print("Допустимые типы данных:", valid_data_types)
Код: Выделить всё
# Extract metadata columns and data types from filter_table
metadata_columns_df = filter_table.select('dwh_column_name', 'dwh_data_type')
metadata_data_dict = {row['dwh_column_name']: row['dwh_data_type'] for row in metadata_columns_df.collect()}
# List to hold columns that need type casting
columns_to_cast = []
# Define a mapping from string data type names to Spark SQL types
type_mapping = {
'string': StringType(),
'int': IntegerType(),
'bigint': LongType(),
'float': FloatType(),
'double': DoubleType(),
'boolean': BooleanType(),
# Add other mappings as necessary
}
# Identify columns with mismatched data types
for column, expected_type in metadata_data_dict.items():
if column in valid_data_types and expected_type != valid_data_types[column]:
# Skip 'phone_format_code' column
if expected_type == 'phone_format_code':
continue
# Get the Spark SQL type from the type_mapping dictionary
spark_type = type_mapping.get(expected_type.lower())
if spark_type is not None:
columns_to_cast.append((column, spark_type))
else:
print(f"Unsupported type => {column}: {expected_type}")
# Apply type casting for identified columns
for column, spark_type in columns_to_cast:
try:
valid_range_df = valid_range_df.withColumn(column, col(column).cast(spark_type))
print(f"Column '{column}' type changed to '{spark_type.simpleString()}'.")
except Exception as e:
print(f"Error casting column '{column}': {e}")
# Print schema of the new DataFrame
valid_range_df.printSchema()
print("The columns to cast are:", columns_to_cast)
return valid_range_df
new_valid_range_df = type_check(valid_range_df, filter_table)
Подробнее здесь: https://stackoverflow.com/questions/787 ... g-the-data