大数据处理:Spark与Flink实战
8 阅读
预计 3 分钟
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擅长流处理,两者互补。
0 条评论 欢迎参与讨论