format
This commit is contained in:
parent
6baccbf23d
commit
5303741624
@ -1,4 +1,4 @@
|
|||||||
from pyspark.sql.types import IntegerType, FloatType
|
from pyspark.sql.types import FloatType
|
||||||
|
|
||||||
from client import spark_app
|
from client import spark_app
|
||||||
from data.data import data_now_hdfs, mysql_jdbc_uri, mysql_jdbc_prop
|
from data.data import data_now_hdfs, mysql_jdbc_uri, mysql_jdbc_prop
|
||||||
|
@ -1,4 +1,5 @@
|
|||||||
import os
|
import os
|
||||||
|
|
||||||
os.environ["JAVA_HOME"] = "/opt/modules/jdk1.8.0_212"
|
os.environ["JAVA_HOME"] = "/opt/modules/jdk1.8.0_212"
|
||||||
os.environ["YARN_CONF_DIR"] = "/usr/hdp/2.6.3.0-235/hadoop-yarn/etc/hadoop/"
|
os.environ["YARN_CONF_DIR"] = "/usr/hdp/2.6.3.0-235/hadoop-yarn/etc/hadoop/"
|
||||||
|
|
||||||
|
@ -1,5 +1,5 @@
|
|||||||
from pyspark.sql.types import DecimalType
|
|
||||||
from pyspark.sql import functions
|
from pyspark.sql import functions
|
||||||
|
from pyspark.sql.types import DecimalType
|
||||||
|
|
||||||
from client import spark_app
|
from client import spark_app
|
||||||
from data.data import data_hour_hdfs, mysql_jdbc_uri, mysql_jdbc_prop
|
from data.data import data_hour_hdfs, mysql_jdbc_uri, mysql_jdbc_prop
|
||||||
|
@ -1,12 +1,12 @@
|
|||||||
from pyspark.sql.types import IntegerType
|
|
||||||
from pyspark.sql import functions
|
from pyspark.sql import functions
|
||||||
|
from pyspark.sql.types import IntegerType
|
||||||
|
|
||||||
from client import spark_app
|
from client import spark_app
|
||||||
from data.data import data_day_hdfs, mysql_jdbc_uri, mysql_jdbc_prop
|
from data.data import data_day_hdfs, mysql_jdbc_uri, mysql_jdbc_prop
|
||||||
|
|
||||||
|
|
||||||
def tem_2024_analyse():
|
def tem_2024_analyse():
|
||||||
print("计算重庆市各个区县 2024 年 每天的 最高气温,最低气温,平均气温")
|
print("计算重庆市各个区县 2024 年 每月的 最高气温,最低气温,平均气温")
|
||||||
app = spark_app("tem_2024_analyse")
|
app = spark_app("tem_2024_analyse")
|
||||||
print("创建应用完成,开始读取数据")
|
print("创建应用完成,开始读取数据")
|
||||||
df = app.read.csv(data_day_hdfs, header=True)
|
df = app.read.csv(data_day_hdfs, header=True)
|
||||||
@ -32,7 +32,7 @@ def tem_2024_analyse():
|
|||||||
df_rain_data.cache()
|
df_rain_data.cache()
|
||||||
print("处理完成,保存数据到数据库")
|
print("处理完成,保存数据到数据库")
|
||||||
df_rain_data.coalesce(1).write.jdbc(mysql_jdbc_uri, "tem_2024_analyse", "ignore", mysql_jdbc_prop)
|
df_rain_data.coalesce(1).write.jdbc(mysql_jdbc_uri, "tem_2024_analyse", "ignore", mysql_jdbc_prop)
|
||||||
print("各个区县 2024 年 每天的 最高气温,最低气温,平均气温计算完毕")
|
print("各个区县 2024 年 每月的 最高气温,最低气温,平均气温计算完毕")
|
||||||
return df_rain_data.head(20)
|
return df_rain_data.head(20)
|
||||||
|
|
||||||
|
|
||||||
|
Loading…
Reference in New Issue
Block a user