Spark和Flink是大数据处理的两大引擎。 ## Spark ```python from pyspark.sql import SparkSession, functions as F spark = SparkSession.builder.appName("Processing").getOrCreate() df = spark.read.parquet("s3://bucket/data/") result = df.filter(F.col("age") > 25).groupBy("city").agg(F.avg("salary")) ``` ## Flink流处理 ```python from pyflink.table import StreamTableEnvironment t_env.execute_sql(""" SELECT window_start, SUM(amount) FROM TABLE(TUMBLE(TABLE orders, DESCRIPTOR(order_time), INTERVAL '1' HOUR)) GROUP BY window_start """) ``` ## 选型建议 | 场景 | 推荐 | |------|------| | 批处理 | Spark | | 流处理 | Flink | | 交互式查询 | Spark SQL | Spark擅长批处理,Flink擅长流处理,两者互补。