from pyspark.sql import SparkSessionfrom pyspark.ml.feature import VectorAssemblerfrom pyspark.ml.clustering import KMeansdef get_top_n_products(spark, n=10): df = spark.read.csv("hdfs://.../taobao_cosmetics.csv", header=True, inferSchema=True) df.createOrReplaceTempView("products") top_n_df = spark.sql(f"SELECT title, salesnum FROM products ORDER BY cast(salesnum as int) DESC LIMIT {n}") return top_n_df.toPandas().to_dict('records')def get_province_distribution(spark): df = spark.read.csv("hdfs://.../taobao_cosmetics.csv", header=True, inferSchema=True) df.createOrReplaceTempView("products") province_df = spark.sql("SELECT province, SUM(cast(salesnum as int)) as total_sales FROM products GROUP BY province") return province_df.toPandas().to_dict('records')def kmeans_product_clustering(spark, k=4): df = spark.read.csv("hdfs://.../taobao_cosmetics.csv", header=True, inferSchema=True) df = df.na.drop(subset=["jiage", "salesnum"]) assembler = VectorAssembler(inputCols=["jiage", "salesnum"], outputCol="features") df_with_features = assembler.transform(df) kmeans = KMeans(featuresCol="features", k=k) model = kmeans.fit(df_with_features) clustered_df = model.transform(df_with_features) result_df = clustered_df.select("title", "jiage", "salesnum", "prediction") return result_df.toPandas().to_dict('records')