结构化数据处理是数据分析和机器学习的基础,主要涉及对结构化数据(如表格数据)的清洗、转换、分析和建模。以下是结构化数据处理的核心流程和技术,结合 Python 和 Spark 示例说明:

一、结构化数据概述

结构化数据通常以二维表格形式存在,具有明确的列名和数据类型,常见于关系型数据库、CSV 文件、Excel 表格等。处理工具包括:

  • Python:Pandas、NumPy、SciPy
  • 分布式计算:Spark SQL、Hive
  • 数据库:SQL(MySQL、PostgreSQL 等)

二、数据读取与探查

1. Python 读取数据

python

运行

import pandas as pd

# 读取CSV文件
df = pd.read_csv('data.csv')

# 读取Excel文件
df = pd.read_excel('data.xlsx')

# 读取数据库
from sqlalchemy import create_engine
engine = create_engine('postgresql://user:password@host:port/dbname')
df = pd.read_sql('SELECT * FROM table', engine)
2. 数据探查

python

运行

# 查看基本信息
print('数据基本信息:')
df.info()

# 查看数据集行数和列数
rows, columns = df.shape

if rows < 1000:
    # 小数据集(行数少于1000)查看全量数据信息
    print('数据全部内容信息:')
    print(df.to_csv(sep='\t', na_rep='nan'))
else:
    # 大数据集查看数据前几行信息
    print('数据前几行内容信息:')
    print(df.head().to_csv(sep='\t', na_rep='nan'))

# 查看数据集行数和列数
rows, columns = df.shape

if rows < 1000:
    # 小数据集(行数少于1000)查看全量数据统计信息
    print('数据全部内容统计信息:')
    print(df.describe(include='all', percentiles=[.25, .5, .75]).to_csv(sep='\t', na_rep='nan'))
else:
    # 大数据集查看数据前几行统计信息
    print('数据前几行内容统计信息:')
    print(df.head().describe(include='all', percentiles=[.25, .5, .75]).to_csv(sep='\t', na_rep='nan'))

三、数据清洗

1. 缺失值处理

python

运行

# 检测缺失值
print('数据缺失值统计:')
print(df.isnull().sum())

# 删除缺失值
df_dropna = df.dropna()

# 填充缺失值
df_filled = df.fillna({
    'numeric_col': df['numeric_col'].mean(),  # 数值列用均值填充
    'category_col': 'unknown'  # 分类列用默认值填充
})
2. 异常值处理

python

运行

# 基于Z-score检测异常值
from scipy import stats

z_scores = np.abs(stats.zscore(df['numeric_col']))
threshold = 3
df_clean = df[z_scores < threshold]  # 删除超过3个标准差的值

# 基于IQR检测异常值
Q1 = df['numeric_col'].quantile(0.25)
Q3 = df['numeric_col'].quantile(0.75)
IQR = Q3 - Q1
lower_bound = Q1 - 1.5 * IQR
upper_bound = Q3 + 1.5 * IQR
df_clean = df[(df['numeric_col'] >= lower_bound) & (df['numeric_col'] <= upper_bound)]
3. 重复值处理

python

运行

# 检测重复值
print('数据重复值统计:')
print(df.duplicated().sum())

# 删除重复值
df_clean = df.drop_duplicates()

四、特征工程

1. 数据类型转换

python

运行

# 转换为数值类型
df['numeric_col'] = pd.to_numeric(df['numeric_col'], errors='coerce')

# 转换为日期类型
df['date_col'] = pd.to_datetime(df['date_col'])

# 提取日期特征
df['year'] = df['date_col'].dt.year
df['month'] = df['date_col'].dt.month
df['day_of_week'] = df['date_col'].dt.dayofweek
2. 分类特征编码

python

运行

# 独热编码
df_encoded = pd.get_dummies(df, columns=['category_col'])

# 标签编码
from sklearn.preprocessing import LabelEncoder

le = LabelEncoder()
df['category_col_encoded'] = le.fit_transform(df['category_col'])
3. 数值特征缩放

python

运行

from sklearn.preprocessing import MinMaxScaler, StandardScaler

# Min-Max缩放
scaler = MinMaxScaler()
df['scaled_col'] = scaler.fit_transform(df[['numeric_col']])

# 标准化
scaler = StandardScaler()
df['standardized_col'] = scaler.fit_transform(df[['numeric_col']])
4. 特征组合与衍生

python

运行

# 特征组合
df['total'] = df['price'] * df['quantity']

# 分箱操作
df['age_group'] = pd.cut(df['age'], bins=[0, 18, 30, 50, 100], labels=['child', 'young', 'adult', 'senior'])

五、数据分析与聚合

1. 基本统计分析

python

运行

# 描述性统计
print(df.describe())

# 相关性分析
print(df.corr())

# 分组聚合
grouped = df.groupby('category_col').agg({
    'numeric_col': ['mean', 'sum', 'std'],
    'id': 'count'
})
2. SQL 查询(使用 Pandas)

python

运行

from pandasql import sqldf

pysqldf = lambda q: sqldf(q, globals())

# 执行SQL查询
result = pysqldf("SELECT category_col, AVG(numeric_col) FROM df GROUP BY category_col")

六、Spark 处理结构化数据

对于大规模数据,可使用 Spark SQL 处理结构化数据:

python

运行

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, avg, sum, count, when

# 创建SparkSession
spark = SparkSession.builder \
    .appName("StructuredDataProcessing") \
    .getOrCreate()

# 读取数据
df_spark = spark.read.csv("data.csv", header=True, inferSchema=True)

# 数据清洗
df_spark = df_spark.na.fill(0, subset=["numeric_col"])  # 填充缺失值
df_spark = df_spark.filter(col("numeric_col") > 0)  # 过滤异常值

# 特征工程
df_spark = df_spark.withColumn("total", col("price") * col("quantity"))  # 特征组合
df_spark = df_spark.withColumn("age_group",  # 分箱操作
    when(col("age") <= 18, "child")
    .when(col("age") <= 30, "young")
    .when(col("age") <= 50, "adult")
    .otherwise("senior")
)

# 分组聚合
result = df_spark.groupBy("category_col") \
    .agg(
        avg("numeric_col").alias("avg_value"),
        sum("total").alias("total_sum"),
        count("*").alias("count")
    )

# 执行SQL查询
df_spark.createOrReplaceTempView("data")
result = spark.sql("SELECT category_col, AVG(numeric_col) as avg_value FROM data GROUP BY category_col")

# 保存结果
result.write.csv("output.csv", header=True)

七、数据可视化

python

运行

import matplotlib.pyplot as plt
import seaborn as sns

# 设置中文字体
plt.rcParams["font.family"] = ["SimHei", "WenQuanYi Micro Hei", "Heiti TC"]

# 柱状图
plt.figure(figsize=(10, 6))
sns.barplot(x='category_col', y='numeric_col', data=df)
plt.title('分类变量与数值变量关系')
plt.show()

# 散点图
plt.figure(figsize=(10, 6))
sns.scatterplot(x='x_col', y='y_col', hue='category_col', data=df)
plt.title('散点图')
plt.show()

# 箱线图
plt.figure(figsize=(10, 6))
sns.boxplot(x='category_col', y='numeric_col', data=df)
plt.title('箱线图')
plt.show()

# 热力图(相关性)
plt.figure(figsize=(10, 8))
sns.heatmap(df.corr(), annot=True, cmap='coolwarm')
plt.title('相关性热力图')
plt.show()

八、最佳实践

  1. 数据质量优先:确保数据清洗彻底,避免垃圾进垃圾出。
  2. 模块化处理:将数据处理流程拆分为独立函数,提高可维护性。
  3. 使用向量化操作:Pandas 和 Spark 的向量化操作比循环效率高得多。
  4. 分布式计算:处理大规模数据时,优先使用 Spark 等分布式框架。
  5. 监控与验证:在关键处理步骤添加数据验证逻辑,确保结果符合预期。

结构化数据处理是数据科学的基础技能,通过熟练掌握上述工具和技术,可以高效地完成从数据清洗到分析建模的全流程。

Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐