首次提交

This commit is contained in:
crystal20277 2022-08-29 10:04:10 +08:00
parent 666c51e9ab
commit d96040bc47
88 changed files with 202790 additions and 2 deletions

3
.idea/.gitignore vendored Normal file
View File

@ -0,0 +1,3 @@
# Default ignored files
/shelf/
/workspace.xml

12
.idea/aiforge_python.iml Normal file
View File

@ -0,0 +1,12 @@
<?xml version="1.0" encoding="UTF-8"?>
<module type="PYTHON_MODULE" version="4">
<component name="NewModuleRootManager">
<content url="file://$MODULE_DIR$" />
<orderEntry type="jdk" jdkName="Python 3.8 (aiforge)" jdkType="Python SDK" />
<orderEntry type="sourceFolder" forTests="false" />
</component>
<component name="PyDocumentationSettings">
<option name="format" value="PLAIN" />
<option name="myDocStringFormat" value="Plain" />
</component>
</module>

View File

@ -0,0 +1,12 @@
<component name="InspectionProjectProfileManager">
<profile version="1.0">
<option name="myName" value="Project Default" />
<inspection_tool class="PyUnresolvedReferencesInspection" enabled="true" level="WARNING" enabled_by_default="true">
<option name="ignoredIdentifiers">
<list>
<option value="project_evaluation_analysis.generate_history_data" />
</list>
</option>
</inspection_tool>
</profile>
</component>

View File

@ -0,0 +1,6 @@
<component name="InspectionProjectProfileManager">
<settings>
<option name="USE_PROJECT_PROFILE" value="false" />
<version value="1.0" />
</settings>
</component>

4
.idea/misc.xml Normal file
View File

@ -0,0 +1,4 @@
<?xml version="1.0" encoding="UTF-8"?>
<project version="4">
<component name="ProjectRootManager" version="2" project-jdk-name="Python 3.8 (aiforge)" project-jdk-type="Python SDK" />
</project>

8
.idea/modules.xml Normal file
View File

@ -0,0 +1,8 @@
<?xml version="1.0" encoding="UTF-8"?>
<project version="4">
<component name="ProjectModuleManager">
<modules>
<module fileurl="file://$PROJECT_DIR$/.idea/aiforge_python.iml" filepath="$PROJECT_DIR$/.idea/aiforge_python.iml" />
</modules>
</component>
</project>

6
.idea/vcs.xml Normal file
View File

@ -0,0 +1,6 @@
<?xml version="1.0" encoding="UTF-8"?>
<project version="4">
<component name="VcsDirectoryMappings">
<mapping directory="$PROJECT_DIR$" vcs="Git" />
</component>
</project>

View File

@ -1,2 +0,0 @@
# aiforge_python

18
config.py Normal file
View File

@ -0,0 +1,18 @@
import os
from datetime import datetime
from utils import create_logger
# 项目根目录
root_path = os.path.abspath(os.path.dirname(__file__))
# 创建日志对象
logger = create_logger(root_path + '/logs/' + str(datetime.date(datetime.now())) + '.log')
# PostgreSQL连接配置
host = '118.31.13.117'
user = 'postgres'
passwd = 'Edu_123123pos'
port = 55432
database = 'opendata'

View File

@ -0,0 +1,36 @@
"""
企业行业地位企业业务方向
"""
import warnings
import pandas as pd
from pandas.io import sql
from datetime import datetime
from sqlalchemy import create_engine
from config import host, port, user, passwd, database
warnings.filterwarnings("ignore")
def generate_org_industry_status_data():
engine = create_engine(f'postgresql+psycopg2://{user}:{passwd}@{host}:{port}/{database}', pool_recycle=3600)
conn = engine.connect()
truncate_sql = f"truncate ai_repo_topic_detail"
sql.execute(truncate_sql, conn)
sql_str = f"""
select
t1.id repo_id, t1.num_watches, t1.num_stars, t1.num_forks, t1.clone_cnt,
t1.num_watches + t1.num_stars + t1.num_forks + t1.clone_cnt num_hots,
t3.id topic_id, t3.name topic_name
from repository t1
left join repo_topic t2 on t1.id = t2.repo_id
left join topic t3 on t2.topic_id = t3."id"
where t3.name is not null
"""
all_df = pd.read_sql(sql_str, conn)
all_df.insert(8, 'created_unix', int(datetime.now().timestamp()))
all_df.to_sql('ai_repo_topic_detail', conn, index=False, if_exists='append')
if __name__ == '__main__':
generate_org_industry_status_data()

View File

@ -0,0 +1,47 @@
"""
企业开源实力排名
"""
import warnings
import pandas as pd
from pandas.io import sql
from datetime import datetime
from sqlalchemy import create_engine
from config import host, port, user, passwd, database
warnings.filterwarnings("ignore")
def generate_org_open_source_rank_data():
engine = create_engine(f'postgresql+psycopg2://{user}:{passwd}@{host}:{port}/{database}', pool_recycle=3600)
conn = engine.connect()
truncate_sql = f"truncate ai_org_open_source_rank_statistics"
sql.execute(truncate_sql, conn)
sql_str = f"""
select
org_id, num_repos, num_forks, num_issues, num_pulls,
num_commit, num_all, rank() over(order by num_all desc) num_rank
from(
select
org_id, count(repo_id) num_repos, sum(num_forks) num_forks,
sum(num_issues) num_issues, sum(num_pulls) num_pulls, sum(num_commit) num_commit,
count(repo_id) + sum(num_forks) + sum(num_issues) + sum(num_pulls) + sum(num_commit) num_all
from(
select
t2.org_id, t2.repo_id, t3.num_forks, t3.num_issues, t3.num_pulls, t3.num_commit
from
opendata.public.user t1
left join team_repo t2 on t1.id = t2.org_id
left join repository t3 on t2.repo_id = t3.id
where t1.type = 1 and t2.org_id is not null
) tt1 GROUP BY tt1.org_id
) ttt1
"""
all_df = pd.read_sql(sql_str, conn)
all_df.insert(8, 'num_all_org', all_df.shape[0])
all_df.insert(9, 'created_unix', int(datetime.now().timestamp()))
all_df.to_sql('ai_org_open_source_rank_statistics', conn, index=False, if_exists='append')
if __name__ == '__main__':
generate_org_open_source_rank_data()

View File

@ -0,0 +1,81 @@
import warnings
from datetime import datetime
from utils import get_conn, get_month_date
warnings.filterwarnings("ignore")
def generate_history_data_by_test():
"""
按月生成历史数据
:return:
"""
conn = get_conn()
curs = conn.cursor()
months = [7, 8]
for month in months:
date_list = get_month_date(2022, month)
for date in date_list:
sql = f"""
INSERT INTO
ai_repository_history
SELECT
nextval('ai_repository_history_id_seq') as id, owner_id, owner_name, lower_name, name, description,
website, original_service_type, original_url, default_branch, creator_id,
(ceil(random() * 10) + num_watches) as num_watches,
(ceil(random() * 10) + num_stars) as num_stars,
(ceil(random() * 10) + num_forks) as num_forks,
(ceil(random() * 10) + num_issues) as num_issues,
(ceil(random() * 10) + num_closed_issues) as num_closed_issues,
(ceil(random() * 10) + num_pulls) as num_pulls,
(ceil(random() * 10) + num_closed_pulls) as num_closed_pulls,
num_milestones, num_closed_milestones,
(ceil(random() * 10) + num_commit) as num_commit, repo_type, is_private, is_empty,
is_archived, is_mirror, status, is_fork, fork_id, is_template, template_id, size,
is_fsck_enabled, close_issues_via_commit_in_any_branch, topics, avatar, contract_address,
balance, block_chain_status,
(ceil(random() * 10) + clone_cnt) as clone_cnt, git_clone_cnt, created_unix, updated_unix,
alias, lower_alias, id as repo_id, '{date}' as date_str, {int(datetime.now().timestamp())}
FROM
repository
"""
curs.execute(sql)
conn.commit()
conn.close()
def generate_history_data():
"""
生成当天数据
:return:
"""
conn = get_conn()
curs = conn.cursor()
date = datetime.now().strftime('%Y-%m-%d')
sql = f"""
INSERT INTO
ai_repository_history
SELECT
nextval('ai_repository_history_id_seq') as id, owner_id, owner_name, lower_name, name, description,
website, original_service_type, original_url, default_branch, creator_id,
num_watches, num_stars, num_forks, num_issues, num_closed_issues, num_pulls, num_closed_pulls,
num_milestones, num_closed_milestones, num_commit, repo_type, is_private, is_empty,
is_archived, is_mirror, status, is_fork, fork_id, is_template, template_id, size,
is_fsck_enabled, close_issues_via_commit_in_any_branch, topics, avatar, contract_address,
balance, block_chain_status, clone_cnt, git_clone_cnt, created_unix, updated_unix,
alias, lower_alias, id as repo_id, '{date}' as date_str, {int(datetime.now().timestamp())}
FROM
repository
"""
curs.execute(sql)
conn.commit()
conn.close()
if __name__ == '__main__':
# 生成测试数据
generate_history_data_by_test()
# 生成正式数据
# generate_history_data()

View File

@ -0,0 +1,67 @@
import warnings
import pandas as pd
from datetime import datetime
from sqlalchemy import create_engine
from utils import get_month_date
from config import host, port, user, passwd, database
warnings.filterwarnings("ignore")
def generate_statistics_data_by_test():
"""
按月生成历史数据
:return:
"""
engine = create_engine(f'postgresql+psycopg2://{user}:{passwd}@{host}:{port}/{database}', pool_recycle=3600)
conn = engine.connect()
months = [7, 8]
for month in months:
for date in get_month_date(2022, month):
generate_statistics_data_single(conn, date)
def generate_statistics_data():
"""
生成当天数据
:return:
"""
engine = create_engine(f'postgresql+psycopg2://{user}:{passwd}@{host}:{port}/{database}', pool_recycle=3600)
conn = engine.connect()
generate_statistics_data_single(conn, datetime.now().strftime('%Y-%m-%d'))
def generate_statistics_data_single(conn, date):
sql = f"""select * from ai_repository_history where date_str = '{date}'"""
all_df = pd.read_sql(sql, conn)
# 发展趋势
tmp_df = all_df[['num_watches', 'num_stars', 'num_forks', 'clone_cnt', 'num_issues', 'num_pulls', 'num_commit']]
all_df['sum'] = tmp_df.sum(axis=1)
all_df['all_avg'] = all_df['sum'].mean()
all_df = all_df.round(2)
db_df = all_df[['repo_id', 'date_str', 'sum', 'all_avg']]
db_df.insert(2, 'type', 'repo_trend')
db_df.insert(3, 'add_time', datetime.now())
db_df = db_df.rename(columns={'sum': 'value', 'all_avg': 'all_value'})
db_df.to_sql('ai_repo_trend_statistics', conn, index=False, if_exists='append')
# 开源潜力
tmp_df = all_df[['num_watches', 'num_stars', 'num_forks', 'clone_cnt']]
all_df['sum'] = tmp_df.sum(axis=1)
all_df['all_avg'] = all_df['sum'].mean()
all_df = all_df.round(2)
db_df = all_df[['repo_id', 'date_str', 'sum', 'all_avg']]
db_df.insert(2, 'type', 'repo_potential')
db_df.insert(3, 'add_time', datetime.now())
db_df = db_df.rename(columns={'sum': 'value', 'all_avg': 'all_value'})
db_df.to_sql('ai_repo_trend_statistics', conn, index=False, if_exists='append')
if __name__ == '__main__':
# 生成测试数据
generate_statistics_data_by_test()
# 生成正式数据
# generate_statistics_data()

View File

@ -0,0 +1,7 @@
from project_evaluation_analysis.generate_history_data import generate_history_data
from project_evaluation_analysis.generate_statistics_data import generate_statistics_data
if __name__ == '__main__':
generate_history_data()
generate_statistics_data()

99
sql/01_2022-08-09.sql Normal file
View File

@ -0,0 +1,99 @@
CREATE TABLE "public"."ai_repository_history" (
"id" int8,
"owner_id" int8,
"owner_name" varchar(255) COLLATE "pg_catalog"."default",
"lower_name" varchar(255) COLLATE "pg_catalog"."default",
"name" varchar(255) COLLATE "pg_catalog"."default",
"description" text COLLATE "pg_catalog"."default",
"website" varchar(2048) COLLATE "pg_catalog"."default",
"original_service_type" int4,
"original_url" varchar(2048) COLLATE "pg_catalog"."default",
"default_branch" varchar(255) COLLATE "pg_catalog"."default",
"creator_id" int8,
"num_watches" int4,
"num_stars" int4,
"num_forks" int4,
"num_issues" int4,
"num_closed_issues" int4,
"num_pulls" int4,
"num_closed_pulls" int4,
"num_milestones" int4,
"num_closed_milestones" int4,
"num_commit" int8,
"repo_type" int4,
"is_private" bool,
"is_empty" bool,
"is_archived" bool,
"is_mirror" bool,
"status" int4,
"is_fork" bool,
"fork_id" int8,
"is_template" bool,
"template_id" int8,
"size" int8,
"is_fsck_enabled" bool,
"close_issues_via_commit_in_any_branch" bool,
"topics" json,
"avatar" varchar(64) COLLATE "pg_catalog"."default",
"contract_address" varchar(255) COLLATE "pg_catalog"."default",
"balance" varchar(255) COLLATE "pg_catalog"."default",
"block_chain_status" int4,
"clone_cnt" int8,
"git_clone_cnt" int8,
"created_unix" int8,
"updated_unix" int8,
"alias" varchar(255) COLLATE "pg_catalog"."default",
"lower_alias" varchar(255) COLLATE "pg_catalog"."default",
"repo_id" int8,
"date_str" varchar(10) COLLATE "pg_catalog"."default"
)
;
ALTER TABLE "public"."ai_repository_history"
OWNER TO "postgres";
CREATE SEQUENCE
ai_repository_history_id_seq
INCREMENT 1
MINVALUE 1
MAXVALUE 9223372036854775807
START WITH 1
CACHE 1;
CREATE SEQUENCE
ai_repo_trend_statistics_seq
INCREMENT 1
MINVALUE 1
MAXVALUE 9223372036854775807
START WITH 1
CACHE 1;
CREATE TABLE "public"."ai_org_open_source_rank_statistics" (
"org_id" int8,
"num_repos" int4,
"num_forks" int4,
"num_issues" int4,
"num_pulls" int4,
"num_commit" int8,
"num_all" int8,
"num_rank" int4,
"num_all_org" int4,
"created_unix" int8
);
CREATE TABLE "public"."ai_repo_topic_detail" (
"repo_id" int8,
"num_watches" int4,
"num_stars" int4,
"num_forks" int4,
"clone_cnt" int8,
"num_hots" int8,
"topic_id" int4,
"topic_name" varchar(255) COLLATE "pg_catalog"."default",
"created_unix" int8
);

View File

@ -0,0 +1,2 @@
# 头歌平台用户画像分析系统

View File

@ -0,0 +1,2 @@
## 这里面存放数据

File diff suppressed because it is too large Load Diff

View File

@ -0,0 +1,87 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from utils import get_before_date
from config import activity_analysis_path
from config import data_before_days
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户活跃度数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text = f"""
SELECT id, user_id, action_type, action_id, created_at, updated_at, ip
FROM user_actions
WHERE created_at >= '{latest_data_date}'
"""
data_df = pd.read_sql(sql_text, con=sql_conn)
data_df.to_csv(activity_analysis_path + 'data/user_actions.csv', index=False, header=True, sep='\t')
logger.info("用户活跃度数据下载完毕,共" + str(data_df.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理用户活跃度数据
"""
data_user_action = pd.read_csv(activity_analysis_path + 'data/user_actions.csv', sep='\t')
data_user_action = data_user_action[['user_id','action_type','created_at']]
# 提取出三种最常用的操作
some_action = ['Attachment','Login','closeNotice', 'clickNotice']
data_user_action = data_user_action.loc[data_user_action['action_type'].isin(some_action)]
# 将日期格式进行转换
data_user_action['created_at'] = pd.to_datetime(data_user_action['created_at'], dayfirst=True)
data_user_action['created_at'] = data_user_action['created_at'].dt.date
return data_user_action
def get_rfm_data():
"""
生成RFM聚类模型训练数据
"""
data = read_action_data()
df_rfm = data[['user_id','action_type','created_at']]
# 给不同操作赋值,来区分不同操作的重要度
class_dict = {
'Attachment': 2,
'Login': 1,
'closeNotice': 4,
'closeNotice': 4
}
df_rfm['action_values'] = df_rfm['action_type'].map(lambda x:x)
df_rfm['action_values'] = df_rfm['action_values'].map(class_dict)
# 计算 R,F,M 值
df_rfm = df_rfm.groupby("user_id").agg({'created_at':'max','user_id':'count','action_values':'sum'})
df_rfm = df_rfm.rename(columns ={'created_at':'Recentdate','user_id':'F','action_values':'M'})
df_rfm["R"] =(df_rfm['Recentdate'].max() - df_rfm['Recentdate'])/np.timedelta64(1,'D')
df_rfm["R"] = df_rfm["R"].astype('str').str.split(" ",expand=True)
df_rfm["R"] = df_rfm["R"].astype("float").astype("int")
df_rfm.drop(columns='Recentdate',inplace=True)
return df_rfm
if __name__ == '__main__':
# 取之前多少天登录的数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1,98 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from utils import get_before_date
from config import activity_analysis_path
from config import data_before_days
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户活跃度数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text = f"""
SELECT
t1.id,
t1.user_id,
t1.action_type,
t1.action_id,
t1.created_at,
t1.updated_at,
t1.ip
FROM
user_actions t1
LEFT JOIN users t2 ON t1.user_id = t2.id
WHERE t1.action_type IN ('Attachment', 'Login', 'closeNotice', 'clickNotice')
AND DATE_FORMAT(t1.created_at, '%Y-%m-%d') >= '{latest_data_date}'
AND DATE_FORMAT(t2.last_login_on, '%Y-%m-%d') >= '{latest_data_date}'
"""
data_df = pd.read_sql(sql_text, con=sql_conn)
data_df.to_csv(activity_analysis_path + 'data/user_actions.csv', index=False, header=True, sep='\t')
logger.info("用户活跃度数据下载完毕,共" + str(data_df.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理用户活跃度数据
"""
data_user_action = pd.read_csv(activity_analysis_path + 'data/user_actions.csv', sep='\t')
data_user_action = data_user_action[['user_id','action_type','created_at']]
# 提取出三种最常用的操作
some_action = ['Attachment','Login','closeNotice', 'clickNotice']
data_user_action = data_user_action.loc[data_user_action['action_type'].isin(some_action)]
# 将日期格式进行转换
data_user_action['created_at'] = pd.to_datetime(data_user_action['created_at'], dayfirst=True)
data_user_action['created_at'] = data_user_action['created_at'].dt.date
return data_user_action
def get_rfm_data():
"""
生成RFM聚类模型训练数据
"""
data = read_action_data()
df_rfm = data[['user_id','action_type','created_at']]
# 给不同操作赋值,来区分不同操作的重要度
class_dict = {
'Attachment': 2,
'Login': 1,
'closeNotice': 4,
'closeNotice': 4
}
df_rfm['action_values'] = df_rfm['action_type'].map(lambda x:x)
df_rfm['action_values'] = df_rfm['action_values'].map(class_dict)
# 计算 R,F,M 值
df_rfm = df_rfm.groupby("user_id").agg({'created_at':'max','user_id':'count','action_values':'sum'})
df_rfm = df_rfm.rename(columns ={'created_at':'Recentdate','user_id':'F','action_values':'M'})
df_rfm["R"] =(df_rfm['Recentdate'].max() - df_rfm['Recentdate'])/np.timedelta64(1,'D')
df_rfm["R"] = df_rfm["R"].astype('str').str.split(" ",expand=True)
df_rfm["R"] = df_rfm["R"].astype("float").astype("int")
df_rfm.drop(columns='Recentdate',inplace=True)
return df_rfm
if __name__ == '__main__':
# 取之前多少天登录的数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1,22 @@
import pickle
from config import logger
from config import test_user_id
from config import activity_analysis_path
logger.info('加载用户活跃度字典')
user_activity_dict = pickle.load(open(activity_analysis_path + 'results/user_activity_dict.pkl', 'rb'))
def user_activity_predict(user_id):
"""
用户活跃度预测
"""
if user_id not in user_activity_dict:
result = -1
else:
result = user_activity_dict[user_id]
return result
if __name__ == '__main__':
result = user_activity_predict(user_id=test_user_id)
print('用户ID', test_user_id, '贡献度:', result)

View File

@ -0,0 +1,2 @@
## 这里存放模型的训练结果

View File

@ -0,0 +1,80 @@
import os
import pandas as pd
import pickle
from sklearn.cluster import KMeans
from config import logger
from config import RANDOM_SEED
from data_process import get_rfm_data
from config import activity_analysis_path
def iflabel(x):
if x == "高高高":
return 5
elif x == "高低高":
return 4
elif x == "高高低":
return 4
elif x == "高低低":
return 3
elif x == "低高高":
return 4
elif x == "低低高":
return 2
elif x == "低高低":
return 2
elif x == "低低低":
return 1
def train():
"""
用户活跃度分析模型训练
"""
logger.info('开始训练用户活跃度分析模型...')
df_rfm = get_rfm_data()
rfm_cluster = df_rfm.copy()
# 根据 R,F,M 值将所有数据分为 8 类
Cluster = df_rfm[["R","F","M"]].values
kmeans_model = KMeans(n_clusters=8, random_state=RANDOM_SEED).fit(Cluster)
rfm_cluster["label"] = kmeans_model.labels_
# 此处计算均值 RFM 模型每个用户的R、F、M值同总体均值进行比较比均值高的就为“高”低的就标明“低”然后通过三个指标的高低来打标签。
df_rfm_mean = rfm_cluster[["label","R","F","M"]].groupby("label").agg("mean")
df_rfm_mean["IFR"] = df_rfm_mean["R"].apply(lambda x: "" if x > df_rfm_mean["R"].mean() else "")
df_rfm_mean["IFF"] = df_rfm_mean["F"].apply(lambda x: "" if x > df_rfm_mean["F"].mean() else "")
df_rfm_mean["IFM"] = df_rfm_mean["M"].apply(lambda x: "" if x > df_rfm_mean["M"].mean() else "")
#方法和上面一样,#将三个指标合并起来生成临时列temp表示RFM综合指标
df_rfm_mean["temp"] = df_rfm_mean["IFR"] +df_rfm_mean["IFF"]+df_rfm_mean["IFM"]
# 对指标进行综合判断
df_rfm_mean["label"] = df_rfm_mean["temp"].apply(lambda x: iflabel(x))
# 此处在使用聚类RFM模型比较的时候使用总体均值进行比较
rfm_cluster["IFR"] = rfm_cluster["R"].apply(lambda x: "" if x > df_rfm_mean["R"].mean() else "")
rfm_cluster["IFF"] = rfm_cluster["F"].apply(lambda x: "" if x > df_rfm_mean["F"].mean() else "")
rfm_cluster["IFM"] = rfm_cluster["M"].apply(lambda x: "" if x > df_rfm_mean["M"].mean() else "")
# 方法和上面一样将三个指标合并起来生成临时列temp表示RFM综合指标
rfm_cluster["temp"] = rfm_cluster["IFR"] + rfm_cluster["IFF"] + rfm_cluster["IFM"]
# 根据定义的打分规则给每个用户进行打分1-5分
rfm_cluster["label"] = rfm_cluster["temp"].apply(lambda x: iflabel(x))
result = rfm_cluster[['label']]
result_csv = pd.DataFrame(list(result.index),index=range(len(list(result.index))),columns = ['user_id'])
result_csv['activity'] = list(result.values.reshape(1,len(list(result.index)))[0])
# 将用户 id 和活跃度分值提取出来并创新建立一个表
result_csv['user_id'] = result_csv['user_id'].astype(int)
result_csv['activity'] = result_csv['activity'].astype(int)
resuts_dict = dict(zip(result_csv['user_id'], result_csv['activity']))
# 保存模型和结果
pickle.dump(kmeans_model, open(activity_analysis_path + 'results/user_activity_model.pkl', 'wb'))
pickle.dump(resuts_dict, open(activity_analysis_path + 'results/user_activity_dict.pkl', 'wb'))
logger.info('用户活跃度分析模型训练完成')
return resuts_dict
if __name__ == '__main__':
train()

View File

@ -0,0 +1,36 @@
import os
from datetime import datetime
from utils import create_logger
# 项目根目录
root_path = os.path.abspath(os.path.dirname(__file__))
# 各子模型目录
activity_analysis_path = root_path + '/activity_analysis/'
contribution_analysis_path = root_path + '/contribution_analysis/'
interests_analysis_path = root_path + '/interests_analysis/'
learning_ability_analysis_path = root_path + '/learning_ability_analysis/'
programming_ability_analysis_path = root_path + '/programming_ability_analysis/'
professional_ability_analysis_path = root_path + '/professional_ability_analysis/'
user_label_analysis_path = root_path + '/user_label_analysis/'
# 随机数种子
RANDOM_SEED = 42
# 创建日志对象
logger = create_logger(root_path + '/logs/' + str(datetime.date(datetime.now())) + '.log')
# mysql连接配置
mysql_host = "rm-bp13v5020p7828r5rso.mysql.rds.aliyuncs.com"
mysql_user = "testeducoder"
mysql_passwd = "TEST@123"
mysql_port = 3306
mysql_database = "preeducoderweb"
# 之前登录的天数
before_login_days = 30
test_user_id = 201
data_before_days = 10

View File

@ -0,0 +1 @@
## 这里面存放数据

View File

@ -0,0 +1 @@
id journalized_id journalized_type user_id notes created_on private_notes parent_id comments_count reply_id
1 id journalized_id journalized_type user_id notes created_on private_notes parent_id comments_count reply_id

View File

@ -0,0 +1,2 @@
id jour_id jour_type user_id notes status reply_id created_on updated_on
112689 432151 HomeworkCommon 21338 11111 1 2022-04-12 14:46:19 2022-04-12 14:46:19
1 id jour_id jour_type user_id notes status reply_id created_on updated_on
2 112689 432151 HomeworkCommon 21338 11111 1 2022-04-12 14:46:19 2022-04-12 14:46:19

View File

@ -0,0 +1,8 @@
user_id notes created_on
15583 讨论帖 2022-04-12 13:53:39
21338 2022-04-12 13:55:10
21338 2022-04-12 18:35:30
21338 2022-04-12 19:14:20
21338 2022-04-13 11:46:06
21338 2022-04-13 14:11:00
21338 2022-04-13 15:11:00
1 user_id notes created_on
2 15583 讨论帖 2022-04-12 13:53:39
3 21338 2022-04-12 13:55:10
4 21338 2022-04-12 18:35:30
5 21338 2022-04-12 19:14:20
6 21338 2022-04-13 11:46:06
7 21338 2022-04-13 14:11:00
8 21338 2022-04-13 15:11:00

View File

@ -0,0 +1,99 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from utils import get_before_date
from config import contribution_analysis_path
from config import data_before_days
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户贡献度数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text1 = f"""
SELECT id, journalized_id, journalized_type, user_id, notes,created_on,private_notes,parent_id,comments_count,reply_id
FROM journals
"""
data_journals = pd.read_sql(sql_text1, con=sql_conn)
data_journals.to_csv(contribution_analysis_path + 'data/journals.csv', index=False, header=True, sep='\t')
logger.info("用户贡献度数据1数据下载完毕" + str(data_journals.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text2 = f"""
SELECT id,jour_id,jour_type,user_id,notes,status,reply_id,created_on,updated_on
FROM journals_for_messages
"""
data_journals_for_messages = pd.read_sql(sql_text2, con=sql_conn)
data_journals_for_messages.to_csv(contribution_analysis_path + 'data/journals_for_messages.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据2下载完毕" + str(data_journals_for_messages.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理用户贡献度度数据
"""
data_journals = pd.read_csv(contribution_analysis_path + 'data/journals.csv', sep='\t')
data_journals = data_journals[['user_id','notes','created_on']]
data_journals_for_message = pd.read_csv(contribution_analysis_path + 'data/journals_for_messages.csv', sep='\t')
data_journals_for_message = data_journals_for_message[['user_id','notes','created_on']]
# 连接两个表
data_concat = pd.concat([data_journals, data_journals_for_message], ignore_index=True)
# 清洗数据,将'\n','\r',' ',全是数字或这字母的评论删去
data_concat = data_concat.dropna()
data_concat['notes'] = data_concat['notes'].str.replace('\n', '')
data_concat['notes'] = data_concat['notes'].str.replace('\r', '')
data_concat['notes'] = data_concat['notes'].str.replace(' ', '')
data_concat['notes'] = data_concat['notes'].str.strip()
data_concat = data_concat[data_concat['notes'].str.isdigit() != True]
data_concat = data_concat[data_concat['notes'].str.isalpha() != True]
# 转换日期格式
data_concat['created_on'] = pd.to_datetime(data_concat['created_on'], dayfirst=True)
data_concat['created_on'] = data_concat['created_on'].dt.date
data_concat = data_concat.dropna()
return data_concat
def get_rfm_data():
"""
生成RFM聚类模型训练数据
"""
data_all = read_action_data()
data_all['value'] = data_all['notes'].str.len()
data_all['value'] = data_all['value'].apply(lambda x: 4 if x > 20 else 2)
# 计算 R,F,M 值
df_rfm = data_all[['user_id','value','created_on']]
df_rfm = df_rfm.groupby("user_id").agg({'created_on': 'max', 'user_id': 'count', 'value': 'sum'})
df_rfm = df_rfm.rename(columns={'created_on': 'Recentdate', 'user_id': 'F', 'value': 'M'})
df_rfm["R"] = (df_rfm['Recentdate'].max() - df_rfm['Recentdate']) / np.timedelta64(1, 'D')
df_rfm["R"] = df_rfm["R"].astype('str').str.split(" ", expand=True)
df_rfm["R"] = df_rfm["R"].astype("float").astype("int")
df_rfm.drop(columns='Recentdate', inplace=True)
return df_rfm
if __name__ == '__main__':
# 取数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1,159 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from utils import get_before_date
from config import contribution_analysis_path
from config import data_before_days
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户贡献度数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text1 = f"""
SELECT
t1.id,
t1.journalized_id,
t1.journalized_type,
t1.user_id,
t1.notes,
t1.created_on,
t1.private_notes,
t1.parent_id,
t1.comments_count,
t1.reply_id
FROM
journals t1
LEFT JOIN users t2 ON t1.user_id = t2.id
WHERE
DATE_FORMAT(t2.last_login_on, '%Y-%m-%d') >= '{latest_data_date}'
AND DATE_FORMAT(t1.created_on, '%Y-%m-%d') >= '{latest_data_date}'
"""
data_journals = pd.read_sql(sql_text1, con=sql_conn)
data_journals.to_csv(contribution_analysis_path + 'data/journals.csv', index=False, header=True, sep='\t')
logger.info("用户贡献度数据1数据下载完毕" + str(data_journals.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text2 = f"""
SELECT
t1.id,
t1.jour_id,
t1.jour_type,
t1.user_id,
t1.notes,
t1.status,
t1.reply_id,
t1.created_on,
t1.updated_on
FROM
journals_for_messages t1
left join users t2 on t1.user_id = t2.id
where
DATE_FORMAT(t2.last_login_on, '%Y-%m-%d') >= '{latest_data_date}'
and DATE_FORMAT(t1.created_on, '%Y-%m-%d') >= '{latest_data_date}'
"""
data_journals_for_messages = pd.read_sql(sql_text2, con=sql_conn)
data_journals_for_messages.to_csv(contribution_analysis_path + 'data/journals_for_messages.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据2下载完毕" + str(data_journals_for_messages.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_msg_discusses = f"""
SELECT
t1.author_id user_id,
t1.subject notes,
t1.created_on
FROM
messages t1
LEFT JOIN users t2 ON t1.author_id = t2.id
WHERE
DATE_FORMAT(t2.last_login_on, '%Y-%m-%d') >= '{latest_data_date}'
AND DATE_FORMAT(t1.created_on, '%Y-%m-%d') >= '{latest_data_date}'
UNION
SELECT
t1.user_id,
t1.content notes,
t1.created_at created_on
FROM
discusses t1
LEFT JOIN users t2 ON t1.user_id = t2.id
WHERE
DATE_FORMAT(t2.last_login_on, '%Y-%m-%d') >= '{latest_data_date}'
AND DATE_FORMAT(t1.created_at, '%Y-%m-%d') >= '{latest_data_date}'
"""
data_msg_discusses = pd.read_sql(sql_msg_discusses, con=sql_conn)
data_msg_discusses.to_csv(contribution_analysis_path + 'data/msg_discusses.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据3下载完毕" + str(data_msg_discusses.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理用户贡献度度数据
"""
data_journals = pd.read_csv(contribution_analysis_path + 'data/journals.csv', sep='\t')
data_journals = data_journals[['user_id','notes','created_on']]
data_journals_for_message = pd.read_csv(contribution_analysis_path + 'data/journals_for_messages.csv', sep='\t')
data_journals_for_message = data_journals_for_message[['user_id','notes','created_on']]
# 连接两个表
data_concat = pd.concat([data_journals, data_journals_for_message], ignore_index=True)
# 清洗数据,将'\n','\r',' ',全是数字或这字母的评论删去
data_concat = data_concat.dropna()
data_concat['notes'] = data_concat['notes'].str.replace('\n', '')
data_concat['notes'] = data_concat['notes'].str.replace('\r', '')
data_concat['notes'] = data_concat['notes'].str.replace(' ', '')
data_concat['notes'] = data_concat['notes'].str.strip()
data_concat = data_concat[data_concat['notes'].str.isdigit() != True]
data_concat = data_concat[data_concat['notes'].str.isalpha() != True]
# 转换日期格式
data_concat['created_on'] = pd.to_datetime(data_concat['created_on'], dayfirst=True)
data_concat['created_on'] = data_concat['created_on'].dt.date
data_concat = data_concat.dropna()
return data_concat
def get_rfm_data():
"""
生成RFM聚类模型训练数据
"""
data_all = read_action_data()
data_all['value'] = data_all['notes'].str.len()
data_all['value'] = data_all['value'].apply(lambda x: 4 if x > 20 else 2)
# 计算 R,F,M 值
df_rfm = data_all[['user_id','value','created_on']]
df_rfm = df_rfm.groupby("user_id").agg({'created_on': 'max', 'user_id': 'count', 'value': 'sum'})
df_rfm = df_rfm.rename(columns={'created_on': 'Recentdate', 'user_id': 'F', 'value': 'M'})
df_rfm["R"] = (df_rfm['Recentdate'].max() - df_rfm['Recentdate']) / np.timedelta64(1, 'D')
df_rfm["R"] = df_rfm["R"].astype('str').str.split(" ", expand=True)
df_rfm["R"] = df_rfm["R"].astype("float").astype("int")
df_rfm.drop(columns='Recentdate', inplace=True)
return df_rfm
if __name__ == '__main__':
# 取数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1,22 @@
import pickle
from config import logger
from config import test_user_id
from config import contribution_analysis_path
logger.info('加载用户贡献度度字典')
user_contribution_dict = pickle.load(open(contribution_analysis_path + 'results/user_contribution_dict.pkl', 'rb'))
def user_contribution_predict(user_id):
"""
用户贡献度预测
"""
if user_id not in user_contribution_dict:
result = -1
else:
result = user_contribution_dict[user_id]
return result
if __name__ == '__main__':
result = user_contribution_predict(user_id=test_user_id)
print('用户ID', test_user_id, '活跃度:', result)

View File

@ -0,0 +1,2 @@
## 这里存放模型的训练结果

View File

@ -0,0 +1,79 @@
import os
import pandas as pd
import pickle
from sklearn.cluster import KMeans
from config import logger
from config import RANDOM_SEED
from contribution_analysis.data_process import get_rfm_data
from config import contribution_analysis_path
def iflabel(x):
# 设置用户贡献度打分规则
if x == "高高高":
return 5
elif x == "高低高":
return 4
elif x == "高高低":
return 4
elif x == "高低低":
return 3
elif x == "低高高":
return 4
elif x == "低低高":
return 2
elif x == "低高低":
return 2
elif x == "低低低":
return 1
def train():
"""
用户贡献度分析模型训练
"""
logger.info('开始训练用户贡献度分析模型...')
df_rfm = get_rfm_data()
rfm_cluster = df_rfm.copy()
# 根据R,F,M值进行聚类分为8类
Cluster = df_rfm[["R", "F", "M"]].values
kmeans_model = KMeans(n_clusters=8, random_state=RANDOM_SEED).fit(Cluster)
rfm_cluster["label"] = kmeans_model.labels_
# 先求每一类的R、F、M均值 # 此处计算均值 RFM 模型每个用户的R、F、M值同总体均值进行比较比均值高的就为“高”低的就标明“低”然后通过三个指标的高低来打标签。
df_rfm_mean = rfm_cluster[["label", "R", "F", "M"]].groupby("label").agg("mean")
df_rfm_mean["IFR"] = df_rfm_mean["R"].apply(lambda x: "" if x > df_rfm_mean["R"].mean() else "")
df_rfm_mean["IFF"] = df_rfm_mean["F"].apply(lambda x: "" if x > df_rfm_mean["F"].mean() else "")
df_rfm_mean["IFM"] = df_rfm_mean["M"].apply(lambda x: "" if x > df_rfm_mean["M"].mean() else "")
# 方法和上面一样,#将三个指标合并起来生成临时列temp表示RFM综合指标
df_rfm_mean["temp"] = df_rfm_mean["IFR"] + df_rfm_mean["IFF"] + df_rfm_mean["IFM"]
#将指标进行综合判断
df_rfm_mean["label"] = df_rfm_mean["temp"].apply(lambda x: iflabel(x))
# 此处在使用聚类RFM模型比较的时候使用总体均值进行比较
rfm_cluster["IFR"] = rfm_cluster["R"].apply(lambda x: "" if x > df_rfm_mean["R"].mean() else "")
rfm_cluster["IFF"] = rfm_cluster["F"].apply(lambda x: "" if x > df_rfm_mean["F"].mean() else "")
rfm_cluster["IFM"] = rfm_cluster["M"].apply(lambda x: "" if x > df_rfm_mean["M"].mean() else "")
## 方法和上面一样将三个指标合并起来生成临时列temp表示RFM综合指标
rfm_cluster["temp"] = rfm_cluster["IFR"] + rfm_cluster["IFF"] + rfm_cluster["IFM"]
rfm_cluster["label"] = rfm_cluster["temp"].apply(lambda x: iflabel(x))
# 根据打分规则函数给用户进行打分
result = rfm_cluster[['label']]
# 将用户 id 和贡献度提取出来保存在字典中
result_csv = pd.DataFrame(list(result.index),index=range(len(list(result.index))),columns = ['user_id'])
result_csv['contribution'] = list(result.values.reshape(1,len(list(result.index)))[0])
result_csv['user_id'] = result_csv['user_id'].astype(int)
result_csv['contribution'] = result_csv['contribution'].astype(int)
results_dict = dict(zip(result_csv['user_id'], result_csv['contribution']))
# 保存模型和结果
pickle.dump(kmeans_model, open(contribution_analysis_path + 'results/user_contribution_model.pkl', 'wb'))
pickle.dump(results_dict, open(contribution_analysis_path + 'results/user_contribution_dict.pkl', 'wb'))
logger.info('用户贡献度度分析模型训练完成')
return results_dict
if __name__ == '__main__':
train()

View File

@ -0,0 +1,193 @@
'''
@Compony: EduCoder
@Author: dengzaiyong
@Date: 2022-03-29 15:16:08
@LastEditTime: 2022-03-29 19:37:08
@LastEditors: dengzaiyong
@Description: 用户画像分析模型预测接口
@FilePath: /user_portrait_analysis/flask_app.py
'''
import json
from flask import Flask
from flask_cors import CORS
from flask import request
from activity_analysis.predict import user_activity_predict
from contribution_analysis.predict import user_contribution_predict
from interests_analysis.predict import user_interests_predict
from professional_ability_analysis.predict import user_professional_ability_predict
from programming_ability_analysis.predict import user_programming_ability_predict
from user_label_analysis.predict import user_label_predict
app = Flask(__name__)
CORS(app, resources=r'/*')
@app.route('/user_activity', methods=["POST"])
def get_user_activity():
'''
以RESTful的方式获取用户活跃度分析模型结果
:param user_id: 用户ID
:return 以json格式返回
'''
result = {}
user_id = request.form.get('user_id', type=str, default='')
if user_id.strip() == '':
result = {
"status_code": str('False'),
"error_msg": str('参数错误: 缺少user_id')
}
return json.dumps(result, ensure_ascii=False)
results = user_activity_predict(int(user_id))
result = {
"status_code": str('True'),
"user_id": str(user_id),
"results": str(results),
}
return json.dumps(result, ensure_ascii=False)
@app.route('/user_contribution', methods=["POST"])
def get_user_contribution():
'''
以RESTful的方式获取用户活跃度分析模型结果
:param user_id: 用户ID
:return 以json格式返回
'''
result = {}
user_id = request.form.get('user_id', type=str, default='')
if user_id.strip() == '':
result = {
"status_code": str('False'),
"error_msg": str('参数错误: 缺少user_id')
}
return json.dumps(result, ensure_ascii=False)
results = user_contribution_predict(int(user_id))
result = {
"status_code": str('True'),
"user_id": str(user_id),
"results": str(results),
}
return json.dumps(result, ensure_ascii=False)
@app.route('/user_interests', methods=["POST"])
def get_user_interests():
'''
以RESTful的方式获取用户活跃度分析模型结果
:param user_id: 用户ID
:return 以json格式返回
'''
result = {}
user_id = request.form.get('user_id', type=str, default='')
if user_id.strip() == '':
result = {
"status_code": str('False'),
"error_msg": str('参数错误: 缺少user_id')
}
return json.dumps(result, ensure_ascii=False)
results1,results2 = user_interests_predict(int(user_id))
result = {
"status_code": str('True'),
"user_id": str(user_id),
"results1": str(results1),
"results2": str(results2),
}
return json.dumps(result, ensure_ascii=False)
@app.route('/user_professional_ability', methods=["POST"])
def get_user_professional_ability():
'''
以RESTful的方式获取用户活跃度分析模型结果
:param user_id: 用户ID
:return 以json格式返回
'''
result = {}
user_id = request.form.get('user_id', type=str, default='')
if user_id.strip() == '':
result = {
"status_code": str('False'),
"error_msg": str('参数错误: 缺少user_id')
}
return json.dumps(result, ensure_ascii=False)
results = user_professional_ability_predict(int(user_id))
result = {
"status_code": str('True'),
"user_id": str(user_id),
"results": str(results),
}
return json.dumps(result, ensure_ascii=False)
@app.route('/user_programming_ability', methods=["POST"])
def get_user_programming_ability():
'''
以RESTful的方式获取用户活跃度分析模型结果
:param user_id: 用户ID
:return 以json格式返回
'''
result = {}
user_id = request.form.get('user_id', type=str, default='')
if user_id.strip() == '':
result = {
"status_code": str('False'),
"error_msg": str('参数错误: 缺少user_id')
}
return json.dumps(result, ensure_ascii=False)
results = user_programming_ability_predict(int(user_id))
result = {
"status_code": str('True'),
"user_id": str(user_id),
"results": str(results),
}
return json.dumps(result, ensure_ascii=False)
@app.route('/user_labels', methods=["POST"])
def get_user_labels():
'''
以RESTful的方式获取用户活跃度分析模型结果
:param user_id: 用户ID
:return 以json格式返回
'''
result = {}
user_id = request.form.get('user_id', type=str, default='')
if user_id.strip() == '':
result = {
"status_code": str('False'),
"error_msg": str('参数错误: 缺少user_id')
}
return json.dumps(result, ensure_ascii=False)
r1,r2,r3,r4 = user_label_predict(int(user_id))
result = {
"status_code": str('True'),
"user_id": str(user_id),
"results1": str(r1),
"results2": str(r2),
"results3": str(r3),
"results4": str(r4),
}
return json.dumps(result, ensure_ascii=False)
# python -m flask run
if __name__ == '__main__':
app.run(host='0.0.0.0', port=8088, debug=True, use_reloader=False)

View File

@ -0,0 +1 @@
## 这里面存放数据

View File

@ -0,0 +1,157 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
from utils import get_before_date
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from config import interests_analysis_path
from config import data_before_days
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户爱好数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text1 = f"""
SELECT user_id,repertoire_id
FROM user_interests
"""
data_shixun = pd.read_sql(sql_text1, con=sql_conn)
data_shixun.to_csv(interests_analysis_path + 'data/user_interests.csv', index=False, header=True, sep='\t')
logger.info("用户爱好数据1数据下载完毕" + str(data_shixun.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text2 = f"""
SELECT id,name
FROM repertoires
"""
data_repertoirse = pd.read_sql(sql_text2, con=sql_conn)
data_repertoirse.to_csv(interests_analysis_path + 'data/repertoires.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据2下载完毕" + str(data_repertoirse.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text3 = f"""
SELECT id,user_id,visits
FROM shixuns
"""
data_shixuns = pd.read_sql(sql_text3, con=sql_conn)
data_shixuns.to_csv(interests_analysis_path + 'data/shixuns.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据3下载完毕" + str(data_shixuns.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text4 = f"""
SELECT shixun_id,tag_repertoire_id
FROM shixun_tag_repertoires
"""
data_shixun_tag_repertoires = pd.read_sql(sql_text4, con=sql_conn)
data_shixun_tag_repertoires.to_csv(interests_analysis_path + 'data/shixun_tag_repertoires.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据4下载完毕" + str(data_shixun_tag_repertoires.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text5 = f"""
SELECT id,sub_repertoire_id,name
FROM tag_repertoires
"""
data_tag_repertoires = pd.read_sql(sql_text5, con=sql_conn)
data_tag_repertoires.to_csv(interests_analysis_path + 'data/tag_repertoires.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据5下载完毕" + str(data_tag_repertoires.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text6 = f"""
SELECT id,name,repertoire_id
FROM sub_repertoires
"""
data_sub_repertoires = pd.read_sql(sql_text6, con=sql_conn)
data_sub_repertoires.to_csv(interests_analysis_path + 'data/sub_repertoires.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据6下载完毕" + str(data_sub_repertoires.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理爱好数据
"""
# 读取所有与用户兴趣表格相关的表格
data_user_interests = pd.read_csv(interests_analysis_path + 'data/user_interests.csv', sep='\t')
data_user_interests= data_user_interests[['user_id','repertoire_id']]
data_repertosires = pd.read_csv(interests_analysis_path + 'data/repertoires.csv', sep='\t')
data_repertosires = data_repertosires[['id','name']]
data_repertosires = data_repertosires.rename(columns={'id': 'repertoire_id'})
data_shixuns = pd.read_csv(interests_analysis_path + 'data/shixuns.csv', sep='\t')
data_shixuns = data_shixuns[['id', 'user_id', 'visits']]
data_shixuns = data_shixuns.rename(columns={'id': 'shixun_id'})
data_shixuns = data_shixuns.dropna()
data_shixuns['user_id'] = data_shixuns['user_id'].astype("float").astype("int")
data_shixuns = data_shixuns[data_shixuns.visits > data_shixuns.visits.mean()]
data_shixun_tag_repertoires = pd.read_csv(interests_analysis_path + 'data/shixun_tag_repertoires.csv', sep='\t')
data_shixun_tag_repertoires = data_shixun_tag_repertoires[['shixun_id','tag_repertoire_id']]
data_shixun_tag_repertoires = data_shixun_tag_repertoires.dropna()
data_tag_repertoires = pd.read_csv(interests_analysis_path + 'data/tag_repertoires.csv', sep='\t')
data_tag_repertoires = data_tag_repertoires[['id','sub_repertoire_id','name']]
data_tag_repertoires = data_tag_repertoires.rename(columns={'id': 'tag_repertoire_id'})
data_sub_repertoires = pd.read_csv(interests_analysis_path + 'data/sub_repertoires.csv', sep='\t')
data_sub_repertoires = data_sub_repertoires[['id', 'name', 'repertoire_id']]
data_sub_repertoires = data_sub_repertoires.rename(columns={'id': 'sub_repertoire_id'})
return data_user_interests,data_repertosires,data_shixuns,data_shixun_tag_repertoires,data_tag_repertoires,data_sub_repertoires
def get_hob_data():
"""
生成用户爱好数据
"""
# 进行树形搜索,将用户感兴趣学科的所有相关的词都给搜索出来作为用户的爱好
df_hob_1, df_hob_2, df_hob_3, df_hob_4, df_hob_5, df_hob_6 = read_action_data()
df_merge_1 = pd.merge(df_hob_4, df_hob_5, on='tag_repertoire_id')
df_merge_2 = df_merge_1[['shixun_id', 'sub_repertoire_id', 'name']]
df_merge_3 = pd.merge(df_merge_2, df_hob_6, on='sub_repertoire_id')
df_merge_3 = df_merge_3[['shixun_id', 'name_x', 'name_y', 'repertoire_id']]
df_merge_4 = pd.merge(df_hob_2, df_merge_3, on='repertoire_id')
df_merge_4 = df_merge_4[['shixun_id', 'name', 'name_x', 'name_y']]
df_merge_5 = df_hob_3[['user_id', 'shixun_id']]
df_merge_6 = pd.merge(df_merge_4, df_merge_5, on='shixun_id')
df_merge_6 = df_merge_6[['user_id', 'name', 'name_x', 'name_y']]
# 将所有表格连接后将所有关键词加起来变为一个整个字符串
df_merge_6['hobbies'] = df_merge_6['name'] + ' ' + df_merge_6['name_x'] + ' ' + df_merge_6['name_y']
data = df_merge_6[['user_id', 'hobbies']]
return data, df_hob_1, df_hob_2
if __name__ == '__main__':
# 取数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1,157 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
from utils import get_before_date
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from config import interests_analysis_path
from config import data_before_days
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户爱好数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text1 = f"""
SELECT user_id,repertoire_id
FROM user_interests
"""
data_shixun = pd.read_sql(sql_text1, con=sql_conn)
data_shixun.to_csv(interests_analysis_path + 'data/user_interests.csv', index=False, header=True, sep='\t')
logger.info("用户爱好数据1数据下载完毕" + str(data_shixun.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text2 = f"""
SELECT id,name
FROM repertoires
"""
data_repertoirse = pd.read_sql(sql_text2, con=sql_conn)
data_repertoirse.to_csv(interests_analysis_path + 'data/repertoires.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据2下载完毕" + str(data_repertoirse.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text3 = f"""
SELECT id,user_id,visits
FROM shixuns
"""
data_shixuns = pd.read_sql(sql_text3, con=sql_conn)
data_shixuns.to_csv(interests_analysis_path + 'data/shixuns.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据3下载完毕" + str(data_shixuns.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text4 = f"""
SELECT shixun_id,tag_repertoire_id
FROM shixun_tag_repertoires
"""
data_shixun_tag_repertoires = pd.read_sql(sql_text4, con=sql_conn)
data_shixun_tag_repertoires.to_csv(interests_analysis_path + 'data/shixun_tag_repertoires.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据4下载完毕" + str(data_shixun_tag_repertoires.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text5 = f"""
SELECT id,sub_repertoire_id,name
FROM tag_repertoires
"""
data_tag_repertoires = pd.read_sql(sql_text5, con=sql_conn)
data_tag_repertoires.to_csv(interests_analysis_path + 'data/tag_repertoires.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据5下载完毕" + str(data_tag_repertoires.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text6 = f"""
SELECT id,name,repertoire_id
FROM sub_repertoires
"""
data_sub_repertoires = pd.read_sql(sql_text6, con=sql_conn)
data_sub_repertoires.to_csv(interests_analysis_path + 'data/sub_repertoires.csv', index=False, header=True, sep='\t')
logger.info(
"用户贡献度数据6下载完毕" + str(data_sub_repertoires.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理爱好数据
"""
# 读取所有与用户兴趣表格相关的表格
data_user_interests = pd.read_csv(interests_analysis_path + 'data/user_interests.csv', sep='\t')
data_user_interests= data_user_interests[['user_id','repertoire_id']]
data_repertosires = pd.read_csv(interests_analysis_path + 'data/repertoires.csv', sep='\t')
data_repertosires = data_repertosires[['id','name']]
data_repertosires = data_repertosires.rename(columns={'id': 'repertoire_id'})
data_shixuns = pd.read_csv(interests_analysis_path + 'data/shixuns.csv', sep='\t')
data_shixuns = data_shixuns[['id', 'user_id', 'visits']]
data_shixuns = data_shixuns.rename(columns={'id': 'shixun_id'})
data_shixuns = data_shixuns.dropna()
data_shixuns['user_id'] = data_shixuns['user_id'].astype("float").astype("int")
data_shixuns = data_shixuns[data_shixuns.visits > data_shixuns.visits.mean()]
data_shixun_tag_repertoires = pd.read_csv(interests_analysis_path + 'data/shixun_tag_repertoires.csv', sep='\t')
data_shixun_tag_repertoires = data_shixun_tag_repertoires[['shixun_id','tag_repertoire_id']]
data_shixun_tag_repertoires = data_shixun_tag_repertoires.dropna()
data_tag_repertoires = pd.read_csv(interests_analysis_path + 'data/tag_repertoires.csv', sep='\t')
data_tag_repertoires = data_tag_repertoires[['id','sub_repertoire_id','name']]
data_tag_repertoires = data_tag_repertoires.rename(columns={'id': 'tag_repertoire_id'})
data_sub_repertoires = pd.read_csv(interests_analysis_path + 'data/sub_repertoires.csv', sep='\t')
data_sub_repertoires = data_sub_repertoires[['id', 'name', 'repertoire_id']]
data_sub_repertoires = data_sub_repertoires.rename(columns={'id': 'sub_repertoire_id'})
return data_user_interests,data_repertosires,data_shixuns,data_shixun_tag_repertoires,data_tag_repertoires,data_sub_repertoires
def get_hob_data():
"""
生成用户爱好数据
"""
# 进行树形搜索,将用户感兴趣学科的所有相关的词都给搜索出来作为用户的爱好
df_hob_1, df_hob_2, df_hob_3, df_hob_4, df_hob_5, df_hob_6 = read_action_data()
df_merge_1 = pd.merge(df_hob_4, df_hob_5, on='tag_repertoire_id')
df_merge_2 = df_merge_1[['shixun_id', 'sub_repertoire_id', 'name']]
df_merge_3 = pd.merge(df_merge_2, df_hob_6, on='sub_repertoire_id')
df_merge_3 = df_merge_3[['shixun_id', 'name_x', 'name_y', 'repertoire_id']]
df_merge_4 = pd.merge(df_hob_2, df_merge_3, on='repertoire_id')
df_merge_4 = df_merge_4[['shixun_id', 'name', 'name_x', 'name_y']]
df_merge_5 = df_hob_3[['user_id', 'shixun_id']]
df_merge_6 = pd.merge(df_merge_4, df_merge_5, on='shixun_id')
df_merge_6 = df_merge_6[['user_id', 'name', 'name_x', 'name_y']]
# 将所有表格连接后将所有关键词加起来变为一个整个字符串
df_merge_6['hobbies'] = df_merge_6['name'] + ' ' + df_merge_6['name_x'] + ' ' + df_merge_6['name_y']
data = df_merge_6[['user_id', 'hobbies']]
return data, df_hob_1, df_hob_2
if __name__ == '__main__':
# 取数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1,33 @@
import pickle
from config import logger
from config import test_user_id
from config import interests_analysis_path
logger.info('加载用户爱好字典')
user_interests_dict = pickle.load(open(interests_analysis_path + 'results/user_interests_dict.pkl', 'rb'))
user_possible_hobby_dict = pickle.load(open(interests_analysis_path + 'results/user_possible_hobby_dict.pkl', 'rb'))
def output_hobbies(hobbies):
str = ''
for i in range(len(hobbies)):
str += hobbies[i]+' '
return str
def user_interests_predict(user_id):
"""
用户爱好预测
"""
if user_id not in user_interests_dict:
result1 = list('')
else:
result1 = user_interests_dict[user_id]
if user_id not in user_possible_hobby_dict:
result2 = list('')
else:
result2 = user_possible_hobby_dict[user_id]
return result1,result2
if __name__ == '__main__':
result1,result2 = user_interests_predict(user_id=test_user_id)
print('用户ID', test_user_id, '的爱好有:', output_hobbies(result1))
print('用户ID', test_user_id, '可能的爱好有:', result2)

View File

@ -0,0 +1,2 @@
## 这里面存放模型训练结果

View File

@ -0,0 +1,185 @@
import os
import pandas as pd
import pickle
from sklearn.cluster import KMeans
from config import logger
from config import RANDOM_SEED
from data_process import get_hob_data
from config import interests_analysis_path
"""
协同过滤算法
"""
from abc import ABCMeta, abstractmethod
import numpy as np
from collections import defaultdict
class CF_base(metaclass=ABCMeta):
def __init__(self, k=3):
self.k = k
self.n_user = None
self.n_item = None
@abstractmethod
def init_param(self, data):
pass
@abstractmethod
def cal_prediction(self, *args):
pass
@abstractmethod
def cal_recommendation(self, user_id, data):
pass
def fit(self, data):
# 计算所有用户的推荐物品
self.init_param(data)
all_users = []
for i in range(self.n_user):
all_users.append(self.cal_recommendation(i, data))
return all_users
class CF_knearest(CF_base):
"""
基于物品的K近邻协同过滤推荐算法
"""
def __init__(self, k, criterion='cosine'):
super(CF_knearest, self).__init__(k)
self.criterion = criterion
self.simi_mat = None
return
def init_param(self, data):
# 初始化参数
self.n_user = data.shape[0]
self.n_item = data.shape[1]
self.simi_mat = self.cal_simi_mat(data)
return
def cal_similarity(self, i, j, data):
# 计算物品i和物品j的相似度
items = data[:, [i, j]]
del_inds = np.where(items == 0)[0]
items = np.delete(items, del_inds, axis=0)
if items.size == 0:
similarity = 0
else:
v1 = items[:, 0]
v2 = items[:, 1]
if self.criterion == 'cosine':
if np.std(v1) > 1e-3: # 方差过大,表明用户间评价尺度差别大需要进行调整
v1 = v1 - v1.mean()
if np.std(v2) > 1e-3:
v2 = v2 - v2.mean()
similarity = (v1 @ v2) / np.linalg.norm(v1, 2) / np.linalg.norm(v2, 2)
elif self.criterion == 'pearson':
similarity = np.corrcoef(v1, v2)[0, 1]
else:
raise ValueError('the method is not supported now')
return similarity
def cal_simi_mat(self, data):
# 计算物品间的相似度矩阵
simi_mat = np.ones((self.n_item, self.n_item))
for i in range(self.n_item):
for j in range(i + 1, self.n_item):
simi_mat[i, j] = self.cal_similarity(i, j, data)
simi_mat[j, i] = simi_mat[i, j]
return simi_mat
def cal_prediction(self, user_row, item_ind):
# 计算预推荐物品i对目标活跃用户u的吸引力
purchase_item_inds = np.where(user_row > 0)[0]
rates = user_row[purchase_item_inds]
simi = self.simi_mat[item_ind][purchase_item_inds]
return np.sum(rates * simi) / np.linalg.norm(simi, 1)
def cal_recommendation(self, user_ind, data):
# 计算目标用户的最具吸引力的k个物品list
item_prediction = defaultdict(float)
user_row = data[user_ind]
un_purchase_item_inds = np.where(user_row == 0)[0]
for item_ind in un_purchase_item_inds:
item_prediction[item_ind] = self.cal_prediction(user_row, item_ind)
res = sorted(item_prediction, key=item_prediction.get, reverse=True)
return res[:self.k]
def str2list(x):
x = x.split(' ')
x = list(set(x))
return x
def train():
"""
用户爱好分析模型训练
"""
logger.info('开始训练爱好分析模型...')
# 将用户爱好数据进行整合
data, df_hob_1, df_hob_2 = get_hob_data()
result = data.groupby('user_id').sum()
# 挖掘用户现有的所有喜好
# 将用户数据中重复的进行删除
result['hobbies'] = result['hobbies'].apply(lambda x: str2list(x))
# 挖掘用户可能的喜好
# 将用户在每个大类的使用情况统计出来
df_hob_1_copy = df_hob_1.copy()
df_hob_1_copy['values'] = True
df_hob_cal = df_hob_1_copy.pivot_table(index='user_id', columns='repertoire_id', values='values',
aggfunc='count').fillna(0)
df_hob_cal = df_hob_cal.astype('int')
# 使用协同过滤算法预测用户最有可能对哪个科目感兴趣
user_hob = np.array(df_hob_cal.values)
cf_model = CF_knearest(k=1)
df_hob_cf = cf_model.fit(user_hob)
# 将预测出来的结果转换为 repertoire_id
y_pred = []
for i in range(len(df_hob_cf)):
if len(df_hob_cf[i]) == 0:
y_pred.append(0)
else:
y_pred.append(df_hob_cf[i][0] + 1)
# 重塑表格将爱好 id 提取出来
df_hob_result = pd.DataFrame(list(df_hob_cal.index), index=range(len(list(df_hob_cal.index))), columns=['user_id'])
df_hob_result['repertoire_id'] = y_pred
# 将 repertoire_id 转化为爱好名称,并删除无效用户 id
df_hob = pd.merge(df_hob_result, df_hob_2, on='repertoire_id', how='outer')
df_hob = df_hob.fillna({'name': ' '})
df_hob = df_hob.dropna()
df_hob = df_hob[['user_id', 'name']]
df_hob = df_hob.rename(columns={'name': 'possible_hobby'})
df_hob['user_id'] = df_hob['user_id'].astype(int)
# 将用户 id 和爱好提取出来并保存
result_csv1 = pd.DataFrame(list(result.index), index=range(len(list(result.index))), columns=['user_id'])
result_csv1['hobbies'] = list(result.hobbies)
result_csv1['user_id'] = result_csv1['user_id'].astype(int)
results_dict1 = dict(zip(result_csv1['user_id'], result_csv1['hobbies']))
# 将用户 id 和爱好提取出来并保存
result_csv2 = pd.DataFrame(list(df_hob.index), index=range(len(list(df_hob.index))), columns=['user_id'])
result_csv2['user_id'] = result_csv2['user_id'].astype(int)
result_csv2['possible_hobby'] = df_hob['possible_hobby']
results_dict2 = dict(zip(result_csv2['user_id'], result_csv2['possible_hobby']))
# 保存模型和结果
pickle.dump(results_dict1, open(interests_analysis_path + 'results/user_interests_dict.pkl', 'wb'))
pickle.dump(results_dict2, open(interests_analysis_path + 'results/user_possible_hobby_dict.pkl', 'wb'))
pickle.dump(cf_model, open(interests_analysis_path + 'results/user_possible_hobby_model.pkl', 'wb'))
logger.info('用户爱好分析模型训练完成')
return results_dict1, results_dict2
if __name__ == '__main__':
train()

View File

@ -0,0 +1 @@
## 这里面存放数据

View File

@ -0,0 +1,7 @@
user_id score created_at
455865 0.0 2022-04-13 15:39:29
455868 0.0 2022-04-13 15:39:29
455865 0.0 2022-04-13 15:41:57
455868 0.0 2022-04-13 15:41:57
455864 0.0 2022-04-13 15:42:20
455864 0.0 2022-04-13 15:42:20
1 user_id score created_at
2 455865 0.0 2022-04-13 15:39:29
3 455868 0.0 2022-04-13 15:39:29
4 455865 0.0 2022-04-13 15:41:57
5 455868 0.0 2022-04-13 15:41:57
6 455864 0.0 2022-04-13 15:42:20
7 455864 0.0 2022-04-13 15:42:20

View File

@ -0,0 +1,40 @@
user_id final_score compelete_status created_at
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
455864 0 2022-04-13 15:42:20
1 user_id final_score compelete_status created_at
2 455864 0 2022-04-13 15:42:20
3 455864 0 2022-04-13 15:42:20
4 455864 0 2022-04-13 15:42:20
5 455864 0 2022-04-13 15:42:20
6 455864 0 2022-04-13 15:42:20
7 455864 0 2022-04-13 15:42:20
8 455864 0 2022-04-13 15:42:20
9 455864 0 2022-04-13 15:42:20
10 455864 0 2022-04-13 15:42:20
11 455864 0 2022-04-13 15:42:20
12 455864 0 2022-04-13 15:42:20
13 455864 0 2022-04-13 15:42:20
14 455864 0 2022-04-13 15:42:20
15 455864 0 2022-04-13 15:42:20
16 455864 0 2022-04-13 15:42:20
17 455864 0 2022-04-13 15:42:20
18 455864 0 2022-04-13 15:42:20
19 455864 0 2022-04-13 15:42:20
20 455864 0 2022-04-13 15:42:20
21 455864 0 2022-04-13 15:42:20
22 455864 0 2022-04-13 15:42:20
23 455864 0 2022-04-13 15:42:20
24 455864 0 2022-04-13 15:42:20
25 455864 0 2022-04-13 15:42:20
26 455864 0 2022-04-13 15:42:20
27 455864 0 2022-04-13 15:42:20
28 455864 0 2022-04-13 15:42:20
29 455864 0 2022-04-13 15:42:20
30 455864 0 2022-04-13 15:42:20
31 455864 0 2022-04-13 15:42:20
32 455864 0 2022-04-13 15:42:20
33 455864 0 2022-04-13 15:42:20
34 455864 0 2022-04-13 15:42:20
35 455864 0 2022-04-13 15:42:20
36 455864 0 2022-04-13 15:42:20
37 455864 0 2022-04-13 15:42:20
38 455864 0 2022-04-13 15:42:20
39 455864 0 2022-04-13 15:42:20
40 455864 0 2022-04-13 15:42:20

View File

@ -0,0 +1 @@
user_id is_finished watch_duration
1 user_id is_finished watch_duration

View File

@ -0,0 +1,128 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from utils import get_before_date
from config import learning_ability_analysis_path
from config import data_before_days
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户学习数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text1 = f"""
SELECT user_id,is_finished,watch_duration
FROM watch_course_videos
"""
data_watch_course_videos = pd.read_sql(sql_text1, con=sql_conn)
data_watch_course_videos.to_csv(learning_ability_analysis_path + 'data/watch_course_videos.csv', index=False, header=True, sep='\t')
logger.info("用户学习数据1下载完毕" + str(data_watch_course_videos.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text2 = f"""
SELECT user_id,final_score,compelete_status,created_at
FROM student_works
"""
data_student_works = pd.read_sql(sql_text2, con=sql_conn)
data_student_works.to_csv(learning_ability_analysis_path + 'data/student_works.csv', index=False, header=True, sep='\t')
logger.info(
"用户学习数据2下载完毕" + str(data_student_works.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text3 = f"""
SELECT user_id,score,created_at
FROM exercise_users
"""
data_exercise = pd.read_sql(sql_text3, con=sql_conn)
data_exercise.to_csv(learning_ability_analysis_path + 'data/exercise_users.csv', index=False, header=True, sep='\t')
logger.info(
"用户学习数据3下载完毕" + str(data_exercise.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理爱好数据
"""
data_watch_course_videos = pd.read_csv(learning_ability_analysis_path+ 'data/watch_course_videos.csv', sep='\t')
data_watch_course_videos = data_watch_course_videos[['user_id','is_finished','watch_duration']]
data_student_works = pd.read_csv(learning_ability_analysis_path + 'data/student_works.csv', sep='\t')
data_student_works = data_student_works[['user_id', 'final_score', 'compelete_status']]
data_student_works = data_student_works.fillna(0)
data_exercise = pd.read_csv(learning_ability_analysis_path + 'data/exercise_users.csv', sep='\t')
data_exercise = data_exercise[['user_id','score']]
# 将所有没成绩的 NaN 值填充为 0
data_exercise= data_exercise.fillna(0)
data_exercise['user_id'] = data_exercise['user_id'].astype("float").astype("int")
return data_watch_course_videos,data_student_works,data_exercise
def get_learn_data():
"""
生成学习数据
"""
df_watch, df_work, df_exercise = read_action_data()
# 计算用户观看视频个数,完成观看视频任务个数,观看时常
df_watch_cal = df_watch.groupby('user_id').agg({'is_finished': 'sum', 'user_id': 'count', 'watch_duration': 'sum'})
df_watch_cal = df_watch_cal.rename(columns={'user_id': 'total'})
# 计算用户完成率,并将上述数据重建一个表保存
df_watch_cal['is_finished'] = df_watch_cal['is_finished'] / df_watch_cal['total']
df_watch_result = pd.DataFrame(list(df_watch_cal.index), index=range(len(list(df_watch_cal.index))),
columns=['user_id'])
df_watch_result['is_finished'] = list(df_watch_cal.is_finished.values.reshape(1, len(list(df_watch_cal.index)))[0])
df_watch_result['total'] = list(df_watch_cal.total.values.reshape(1, len(list(df_watch_cal.index)))[0])
df_watch_result['watch_duration'] = list(df_watch_cal.watch_duration.values.reshape(1, len(list(df_watch_cal.index)))[0])
# 计算用户作业个数,作业分数和实际完成作业个数
df_work_cal = df_work.groupby('user_id').agg({'user_id': 'count', 'final_score': 'sum', 'compelete_status': 'sum'})
df_work_cal = df_work_cal.rename(columns={'user_id': 'total_work'})
# 计算用户作业完成率,并将上述数据重建一个表进行保存
df_work_cal['compelete_status'] = df_work_cal['compelete_status'] / df_work_cal['total_work']
df_work_result = pd.DataFrame(list(df_work_cal.index), index=range(len(list(df_work_cal.index))),
columns=['user_id'])
df_work_result['compelete_status'] = list(
df_work_cal.compelete_status.values.reshape(1, len(list(df_work_cal.index)))[0])
df_work_result['total_work'] = list(df_work_cal.total_work.values.reshape(1, len(list(df_work_cal.index)))[0])
df_work_result['final_score'] = list(df_work_cal.final_score.values.reshape(1, len(list(df_work_cal.index)))[0])
# 计算用户考试次数和得分,并将数据重建表保存
df_exercise_cal = df_exercise.groupby('user_id').agg({'user_id': 'count', 'score': 'sum'})
df_exercise_result = pd.DataFrame(list(df_exercise_cal.index), index=range(len(list(df_exercise_cal.index))),
columns=['user_id'])
df_exercise_result['score'] = list(df_exercise_cal.score.values.reshape(1, len(list(df_exercise_cal.index)))[0])
df_exercise_result['user_id'] = df_exercise_result['user_id'].astype("float").astype("int")
# 将上述数据合并,计算实际完成的视频和作业数
df_merge_1 = pd.merge(df_watch_result, df_work_result)
df_merge_2 = pd.merge(df_exercise_result, df_merge_1)
df_merge_2['watch'] = df_merge_2['is_finished'] * df_merge_2['total']
df_merge_2['work'] = df_merge_2['compelete_status'] * df_merge_2['total_work']
data = df_merge_2[['user_id', 'watch', 'work', 'final_score', 'score']]
return data
if __name__ == '__main__':
# 取数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1,153 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from utils import get_before_date
from config import learning_ability_analysis_path
from config import data_before_days
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户学习数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text1 = f"""
SELECT
t1.user_id,
t1.is_finished,
t1.watch_duration
FROM
watch_course_videos t1
LEFT JOIN users t2 ON t1.user_id = t2.id
where
DATE_FORMAT(t2.last_login_on, '%Y-%m-%d') >= '{latest_data_date}'
and DATE_FORMAT(t1.created_at, '%Y-%m-%d') >= '{latest_data_date}'
"""
data_watch_course_videos = pd.read_sql(sql_text1, con=sql_conn)
data_watch_course_videos.to_csv(learning_ability_analysis_path + 'data/watch_course_videos.csv', index=False, header=True, sep='\t')
logger.info("用户学习数据1下载完毕" + str(data_watch_course_videos.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text2 = f"""
SELECT
t1.user_id,
t1.final_score,
t1.compelete_status,
t1.created_at
FROM
student_works t1
LEFT JOIN users t2 ON t1.user_id = t2.id
where
DATE_FORMAT(t2.last_login_on, '%Y-%m-%d') >= '{latest_data_date}'
and DATE_FORMAT(t1.created_at, '%Y-%m-%d') >= '{latest_data_date}'
"""
data_student_works = pd.read_sql(sql_text2, con=sql_conn)
data_student_works.to_csv(learning_ability_analysis_path + 'data/student_works.csv', index=False, header=True, sep='\t')
logger.info(
"用户学习数据2下载完毕" + str(data_student_works.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text3 = f"""
SELECT
t1.user_id,
t1.score,
t1.created_at
FROM
exercise_users t1
LEFT JOIN users t2 ON t1.user_id = t2.id
where
DATE_FORMAT(t2.last_login_on, '%Y-%m-%d') >= '{latest_data_date}'
and DATE_FORMAT(t1.created_at, '%Y-%m-%d') >= '{latest_data_date}'
"""
data_exercise = pd.read_sql(sql_text3, con=sql_conn)
data_exercise.to_csv(learning_ability_analysis_path + 'data/exercise_users.csv', index=False, header=True, sep='\t')
logger.info(
"用户学习数据3下载完毕" + str(data_exercise.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理爱好数据
"""
data_watch_course_videos = pd.read_csv(learning_ability_analysis_path+ 'data/watch_course_videos.csv', sep='\t')
data_watch_course_videos = data_watch_course_videos[['user_id','is_finished','watch_duration']]
data_student_works = pd.read_csv(learning_ability_analysis_path + 'data/student_works.csv', sep='\t')
data_student_works = data_student_works[['user_id', 'final_score', 'compelete_status']]
data_student_works = data_student_works.fillna(0)
data_exercise = pd.read_csv(learning_ability_analysis_path + 'data/exercise_users.csv', sep='\t')
data_exercise = data_exercise[['user_id','score']]
# 将所有没成绩的 NaN 值填充为 0
data_exercise= data_exercise.fillna(0)
data_exercise['user_id'] = data_exercise['user_id'].astype("float").astype("int")
return data_watch_course_videos,data_student_works,data_exercise
def get_learn_data():
"""
生成学习数据
"""
df_watch, df_work, df_exercise = read_action_data()
# 计算用户观看视频个数,完成观看视频任务个数,观看时常
df_watch_cal = df_watch.groupby('user_id').agg({'is_finished': 'sum', 'user_id': 'count', 'watch_duration': 'sum'})
df_watch_cal = df_watch_cal.rename(columns={'user_id': 'total'})
# 计算用户完成率,并将上述数据重建一个表保存
df_watch_cal['is_finished'] = df_watch_cal['is_finished'] / df_watch_cal['total']
df_watch_result = pd.DataFrame(list(df_watch_cal.index), index=range(len(list(df_watch_cal.index))),
columns=['user_id'])
df_watch_result['is_finished'] = list(df_watch_cal.is_finished.values.reshape(1, len(list(df_watch_cal.index)))[0])
df_watch_result['total'] = list(df_watch_cal.total.values.reshape(1, len(list(df_watch_cal.index)))[0])
df_watch_result['watch_duration'] = list(df_watch_cal.watch_duration.values.reshape(1, len(list(df_watch_cal.index)))[0])
# 计算用户作业个数,作业分数和实际完成作业个数
df_work_cal = df_work.groupby('user_id').agg({'user_id': 'count', 'final_score': 'sum', 'compelete_status': 'sum'})
df_work_cal = df_work_cal.rename(columns={'user_id': 'total_work'})
# 计算用户作业完成率,并将上述数据重建一个表进行保存
df_work_cal['compelete_status'] = df_work_cal['compelete_status'] / df_work_cal['total_work']
df_work_result = pd.DataFrame(list(df_work_cal.index), index=range(len(list(df_work_cal.index))),
columns=['user_id'])
df_work_result['compelete_status'] = list(
df_work_cal.compelete_status.values.reshape(1, len(list(df_work_cal.index)))[0])
df_work_result['total_work'] = list(df_work_cal.total_work.values.reshape(1, len(list(df_work_cal.index)))[0])
df_work_result['final_score'] = list(df_work_cal.final_score.values.reshape(1, len(list(df_work_cal.index)))[0])
# 计算用户考试次数和得分,并将数据重建表保存
df_exercise_cal = df_exercise.groupby('user_id').agg({'user_id': 'count', 'score': 'sum'})
df_exercise_result = pd.DataFrame(list(df_exercise_cal.index), index=range(len(list(df_exercise_cal.index))),
columns=['user_id'])
df_exercise_result['score'] = list(df_exercise_cal.score.values.reshape(1, len(list(df_exercise_cal.index)))[0])
df_exercise_result['user_id'] = df_exercise_result['user_id'].astype("float").astype("int")
# 将上述数据合并,计算实际完成的视频和作业数
df_merge_1 = pd.merge(df_watch_result, df_work_result)
df_merge_2 = pd.merge(df_exercise_result, df_merge_1)
df_merge_2['watch'] = df_merge_2['is_finished'] * df_merge_2['total']
df_merge_2['work'] = df_merge_2['compelete_status'] * df_merge_2['total_work']
data = df_merge_2[['user_id', 'watch', 'work', 'final_score', 'score']]
return data
if __name__ == '__main__':
# 取数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1,22 @@
import pickle
from config import logger
from config import test_user_id
from config import learning_ability_analysis_path
logger.info('加载用户学习能力字典')
user_learning_ability_dict = pickle.load(open(learning_ability_analysis_path + 'results/user_learning_ability_dict.pkl', 'rb'))
def user_learning_ability(user_id):
"""
用户学习能力预测
"""
if user_id not in user_learning_ability_dict:
result = -1
else:
result = user_learning_ability_dict[user_id]
return result
if __name__ == '__main__':
result = user_learning_ability(user_id=test_user_id)
print('用户ID', test_user_id, '学习能力:', result)

View File

@ -0,0 +1,2 @@
## 这里面存放模型训练结果

View File

@ -0,0 +1,72 @@
import os
import pandas as pd
import pickle
from sklearn.cluster import KMeans
from config import logger
from config import RANDOM_SEED
from data_process import get_learn_data
from config import learning_ability_analysis_path
def iflabel(x):
if x == 4:
return 5
elif x == 3:
return 4
elif x == 2:
return 3
elif x == 1:
return 2
elif x == 0:
return 1
def train():
"""
用户学习能力模型训练
"""
logger.info
logger.info('开始训练用户学习能力分析模型...')
data = get_learn_data()
# 根据观看视频个数,作业,作业分数,考试分数进行聚类
df_cluster = data[['user_id', 'watch', 'work', 'final_score', 'score']]
kmeans_model = KMeans(n_clusters=8, random_state=RANDOM_SEED).fit(df_cluster)
df_cluster["label"] = kmeans_model.labels_
df_mean = df_cluster[['watch', 'work', 'final_score', 'score', 'label']].groupby("label").agg("mean")
# 这里使用先计算均值再计算分数
# 将每一个指标计算均值超多均值得1分未超过得0分
df_mean["IFWA"] = df_mean["watch"].apply(lambda x: 1 if x > df_mean["watch"].mean() else 0)
df_mean["IFWO"] = df_mean["work"].apply(lambda x: 1 if x > df_mean["work"].mean() else 0)
df_mean["IFF"] = df_mean["final_score"].apply(lambda x: 1 if x > df_mean["final_score"].mean() else 0)
df_mean["IFS"] = df_mean["score"].apply(lambda x: 1 if x > df_mean["score"].mean() else 0)
# 将各项指标分数求和汇总然后根据打分函数进行打分
df_mean["temp"] = df_mean["IFWA"] + df_mean["IFWO"] + df_mean["IFF"] + df_mean["IFS"]
df_mean["label"] = df_mean["temp"].apply(lambda x: iflabel(x))
# 这里使用聚类在进行打分
# 和上面进行相同的操作,最后返回每个用户的学习能力分数
df_cluster["IFWA"] = df_cluster["watch"].apply(lambda x: 1 if x > df_mean["watch"].mean() else 0)
df_cluster["IFWO"] = df_cluster["work"].apply(lambda x: 1 if x > df_mean["work"].mean() else 0)
df_cluster["IFF"] = df_cluster["final_score"].apply(lambda x: 1 if x > df_mean["final_score"].mean() else 0)
df_cluster["IFS"] = df_cluster["score"].apply(lambda x: 1 if x > df_mean["score"].mean() else 0)
df_cluster["temp"] = df_cluster["IFWA"] + df_cluster["IFWO"] + df_cluster["IFF"] + df_cluster["IFS"]
df_cluster["label"] = df_cluster["temp"].apply(lambda x: iflabel(x))
# 将用户 id 和学习能力提取出来保存
result = df_cluster[['label']]
result_csv = pd.DataFrame(list(result.index), index=range(len(list(result.index))), columns=['user_id'])
result_csv['learning_ability'] = list(result.values.reshape(1, len(list(result.index)))[0])
result_csv['user_id'] = result_csv['user_id'].astype(int)
result_csv['learning_ability'] = result_csv['learning_ability'].astype(int)
results_dict = dict(zip(result_csv['user_id'], result_csv['learning_ability']))
# 保存模型和结果
pickle.dump(kmeans_model, open(learning_ability_analysis_path + 'results/user_learning_ability_model.pkl', 'wb'))
pickle.dump(results_dict, open(learning_ability_analysis_path + 'results/user_learning_ability_dict.pkl', 'wb'))
logger.info('用户学习能力分析模型训练完成')
return results_dict
if __name__ == '__main__':
train()

View File

@ -0,0 +1,2 @@
## 这里面存放的是数据

View File

@ -0,0 +1,146 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
from utils import get_before_date
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from config import professional_ability_analysis_path
from config import data_before_days
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户专业知识好数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text1 = f"""
SELECT id,name
FROM repertoires
"""
data_repertoires = pd.read_sql(sql_text1, con=sql_conn)
data_repertoires.to_csv(professional_ability_analysis_path + 'data/repertoires.csv', index=False, header=True, sep='\t')
logger.info("用户专业知识能力数据1数据下载完毕" + str(data_repertoires.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text2 = f"""
SELECT id,shixun_id,user_id,score
FROM challenges
"""
data_challenges = pd.read_sql(sql_text2, con=sql_conn)
data_challenges.to_csv(professional_ability_analysis_path + 'data/challenges.csv', index=False, header=True, sep='\t')
logger.info(
"用户专业知识能力数据2下载完毕" + str(data_challenges.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text3 = f"""
SELECT shixun_id,tag_repertoire_id
FROM shixun_tag_repertoires
"""
data_shixun_tag_repertoires = pd.read_sql(sql_text3, con=sql_conn)
data_shixun_tag_repertoires.to_csv(professional_ability_analysis_path + 'data/shixun_tag_repertoires.csv', index=False, header=True, sep='\t')
logger.info(
"用户专业知识能力数据3下载完毕" + str(data_shixun_tag_repertoires.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text4 = f"""
SELECT id,sub_repertoire_id
FROM tag_repertoires
"""
data_tag_repertoires = pd.read_sql(sql_text4, con=sql_conn)
data_tag_repertoires.to_csv(professional_ability_analysis_path + 'data/tag_repertoires.csv', index=False, header=True, sep='\t')
logger.info(
"用户专业知识能力数据4下载完毕" + str(data_tag_repertoires.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text5 = f"""
SELECT id,repertoire_id
FROM sub_repertoires
"""
data_sub_repertoires = pd.read_sql(sql_text5, con=sql_conn)
data_sub_repertoires.to_csv(professional_ability_analysis_path + 'data/sub_repertoires.csv', index=False, header=True, sep='\t')
logger.info(
"用户专业知识能力数据5下载完毕" + str(data_sub_repertoires.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理用户专业知识能力数据
"""
data_repertoires= pd.read_csv(professional_ability_analysis_path + 'data/repertoires.csv', sep='\t')
data_repertoires = data_repertoires[['id', 'name']]
data_repertoires = data_repertoires.rename(columns={'id': 'repertoire_id'})
data_challenges = pd.read_csv(professional_ability_analysis_path + 'data/challenges.csv', sep='\t')
data_challenges = data_challenges[['user_id', 'id', 'shixun_id', 'score']]
data_challenges = data_challenges.rename(columns={'id': 'challenge_id'})
data_challenges = data_challenges[['user_id', 'shixun_id', 'score']]
data_shixun_tag_repertoires = pd.read_csv(professional_ability_analysis_path + 'data/shixun_tag_repertoires.csv', sep='\t')
data_shixun_tag_repertoires = data_shixun_tag_repertoires[['shixun_id', 'tag_repertoire_id']]
data_shixun_tag_repertoires = data_shixun_tag_repertoires.dropna()
data_tag_repertoires = pd.read_csv(professional_ability_analysis_path + 'data/tag_repertoires.csv', sep='\t')
data_tag_repertoires = data_tag_repertoires[['id', 'sub_repertoire_id']]
data_tag_repertoires = data_tag_repertoires.rename(columns={'id': 'tag_repertoire_id'})
data_sub_repertoires = pd.read_csv(professional_ability_analysis_path + 'data/sub_repertoires.csv', sep='\t')
data_sub_repertoires = data_sub_repertoires[['id', 'repertoire_id']]
data_sub_repertoires = data_sub_repertoires.rename(columns={'id': 'sub_repertoire_id'})
return data_repertoires,data_challenges,data_shixun_tag_repertoires,data_tag_repertoires,data_sub_repertoires
def get_profession_data():
"""
生成用户专业知识能力数据
"""
# 提取用户实训进行搜索,将用户所做实训的大类挑选出来
data_repertoires, data_challenges, data_shixun_tag_repertoires, data_tag_repertoires, data_sub_repertoires = read_action_data()
data_merge_1 = pd.merge(data_shixun_tag_repertoires, data_tag_repertoires, on='tag_repertoire_id')
data_merge_1 = data_merge_1[['shixun_id', 'sub_repertoire_id']]
data_merge_2 = pd.merge(data_merge_1, data_sub_repertoires, on='sub_repertoire_id')
data_merge_2 = data_merge_2[['shixun_id', 'repertoire_id']]
data_merge_3 = pd.merge(data_challenges, data_merge_2, on='shixun_id')
data_merge_3 = data_merge_3[['user_id', 'score', 'repertoire_id']]
cols1 = list(set(list(data_merge_3.repertoire_id.values)))
cols2 = list(data_repertoires.repertoire_id.values)
cols3 = list(set(cols1).union(set(cols2)))
data = data_merge_3.pivot_table(index='user_id', columns='repertoire_id', values='score', aggfunc='sum').fillna(0)
# 根据用户在每个专业类别的得分给用户在每个专业的能力进行打分
for col in list(data.columns):
data.sort_values(by=col, ascending=False)
data[col] = data[col].rank(ascending=1, method='first')
data[col] = data[col] / data[col].max()
q = pd.cut(data[col], 10, labels=[1, 2, 3, 4, 5, 6, 7, 8, 9, 10])
data[col] = q
# 数据库更新时,专业表的名称会增删,此处将删掉的专业列删掉,新增的专业加上去
for i in range(len(cols1)):
cols3.remove(cols1[i])
for i in cols3:
data[i] = 0
return data, data_repertoires
if __name__ == '__main__':
# 取数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1,23 @@
import pickle
from config import logger
from config import test_user_id
from config import activity_analysis_path
from config import professional_ability_analysis_path
logger.info('加载用户专业能力字典')
user_professional_ability_dict = pickle.load(open(professional_ability_analysis_path + 'results/user_professional_ability_dict.pkl', 'rb'))
def user_professional_ability_predict(user_id):
"""
用户专业能力预测
"""
if user_id not in user_professional_ability_dict:
result = -1
else:
result = user_professional_ability_dict[user_id]
return result
if __name__ == '__main__':
result = user_professional_ability_predict(user_id=test_user_id)
print('用户ID', test_user_id, '专业编程能力:', result)

View File

@ -0,0 +1,2 @@
## 这里面存放模型训练结果

View File

@ -0,0 +1,88 @@
import os
import pandas as pd
import pickle
from sklearn.cluster import KMeans
from config import logger
from config import RANDOM_SEED
from data_process import get_profession_data
from config import professional_ability_analysis_path
def iflabel(x):
if x == 8 :
return 5
elif x == 7 :
return 4
elif x == 6:
return 4
elif x == 5:
return 3
elif x == 4:
return 3
elif x == 3:
return 2
elif x == 2:
return 2
elif x == 1:
return 1
elif x == 0:
return 1
def train():
"""
用户专业知识能力分析模型训练
"""
logger.info('开始训练用户专业知识能力分析模型...')
data, data_repertoires = get_profession_data()
cols1 = list(data_repertoires.repertoire_id.values)
# 根据用户在每个专业的得分作为特征进行聚类
df_cluster = data[cols1]
kmeans_model = KMeans(n_clusters=8, random_state=RANDOM_SEED).fit(df_cluster)
df_cluster["label"] = kmeans_model.labels_
# 重建一个表
user_index = list(df_cluster.index)
cols2 = [str(x) for x in cols1]
cols3 = cols2.copy()
cols3.append("label")
df_cluster = pd.DataFrame(df_cluster.values, index=range(len(list(df_cluster.index))), columns=cols3)
df_cluster['user_id'] = user_index
# 计算每类均值并将用户每类的值和均值进行比较大于均值得1分小于均值得0分然后根据打分函数进行打分
df_mean = df_cluster[cols3].groupby("label").agg("mean")
for num in cols2:
df_mean["IF" + num] = df_mean[num].apply(lambda x: 1 if x > df_mean[num].mean() else 0)
df_mean["temp"] = 0
for num in cols2:
df_mean["temp"] += df_mean["IF" + num]
df_mean["label"] = df_mean["temp"].apply(lambda x: iflabel(x))
for num in cols2:
df_cluster["IF" + num] = df_cluster[num].apply(lambda x: 1 if x > df_mean[num].mean() else 0)
df_cluster["temp"] = 0
for num in cols2:
df_cluster["temp"] += df_cluster["IF" + num]
# 将用户 id 和专业能力得分提取出来保存
df_cluster["label"] = df_cluster["temp"].apply(lambda x: iflabel(x))
df_cluster = df_cluster.rename(columns={'label': 'professional_ability'})
result_csv = df_cluster[['user_id', 'professional_ability']]
result_csv['user_id'] = result_csv['user_id'].astype(int)
result_csv['professional_ability'] = result_csv['professional_ability'].astype(int)
results_dict = dict(zip(result_csv['user_id'], result_csv['professional_ability']))
# 保存模型和结果
pickle.dump(kmeans_model,
open(professional_ability_analysis_path + 'results/user_professional_ability_model.pkl', 'wb'))
pickle.dump(results_dict,
open(professional_ability_analysis_path + 'results/user_professional_ability_dict.pkl', 'wb'))
logger.info('用户专业知识能力分析模型训练完成')
return results_dict
if __name__ == '__main__':
train()

View File

@ -0,0 +1,2 @@
## 这里面存放的是数据

View File

@ -0,0 +1,95 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from utils import get_before_date
from config import programming_ability_analysis_path
from config import data_before_days
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户标签数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text1 = f"""
SELECT id, user_id, language
FROM shixuns
"""
data_shixuns = pd.read_sql(sql_text1, con=sql_conn)
data_shixuns.to_csv(programming_ability_analysis_path + 'data/shixuns.csv', index=False, header=True, sep='\t')
logger.info("用户编程数据下载完毕,共" + str(data_shixuns.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理用户编程数据
"""
data_shixuns = pd.read_csv(programming_ability_analysis_path + 'data/shixuns.csv', sep='\t')
data_shixuns = data_shixuns[['id', 'user_id', 'language']]
data_shixuns = data_shixuns.fillna({'language': 'Other'})
data_shixuns = data_shixuns.dropna()
data_shixuns["user_id"] = data_shixuns["user_id"].astype("float").astype("int")
map_list = {
'Java': 'Java',
'MySQL/Java': 'MySQL',
'Python': 'Python',
'C++': 'C/C++',
'Python3.6': 'Python',
'JFinal': 'JFinal',
'Python2.7': 'Python',
'Dynamips': 'Dynamips',
'Ethereum': 'Ethereum',
'Html': 'Html',
'MachineLearning': 'MachineLearning',
'Docker': 'Docker',
'C': 'C/C++',
'MySQL/Python3.6': 'MySQL',
'Verilog': 'Verilog',
'PHP/Web': 'PHP/Web',
'Android': 'Android',
'Golang': 'Golang',
'Hadoop': 'Hadoop',
'Matlab': 'Matlab',
'Shell': 'Shell',
'Git': 'Git',
'Ruby': 'Ruby',
'Perl6': 'Perl6',
'Kotlin': 'Kotlin',
'JavaScript': 'JavaScript',
'Other': 'Other'
}
data_shixuns['language'] = data_shixuns['language'].map(map_list)
return data_shixuns
def get_program_data():
"""
生成用户编程训练数据
"""
data = read_action_data()
data = data[['user_id', 'language']]
# 计算用户在每一类语言实训的总数
data['value'] = 1
data = data.pivot_table(index='user_id', columns='language', values='value', aggfunc='count').fillna(0)
data = data.astype('int')
# 计算用户每类语言的使用占比
data = data.div(data.sum(axis=0), axis=1)
return data
if __name__ == '__main__':
# 取数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1,104 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from utils import get_before_date
from config import programming_ability_analysis_path
from config import data_before_days
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户标签数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text1 = f"""
SELECT
t1.id,
t1.user_id,
t2.language
FROM
myshixuns t1
LEFT JOIN shixuns t2 ON t1.shixun_id = t2.id
LEFT JOIN users t3 ON t1.user_id = t3.id
WHERE
DATE_FORMAT(t3.last_login_on, '%Y-%m-%d') >= '{latest_data_date}'
AND DATE_FORMAT(t1.created_at, '%Y-%m-%d') >= '{latest_data_date}'
"""
data_shixuns = pd.read_sql(sql_text1, con=sql_conn)
data_shixuns.to_csv(programming_ability_analysis_path + 'data/shixuns.csv', index=False, header=True, sep='\t')
logger.info("用户编程数据下载完毕,共" + str(data_shixuns.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理用户编程数据
"""
data_shixuns = pd.read_csv(programming_ability_analysis_path + 'data/shixuns.csv', sep='\t')
data_shixuns = data_shixuns[['id', 'user_id', 'language']]
data_shixuns = data_shixuns.fillna({'language': 'Other'})
data_shixuns = data_shixuns.dropna()
data_shixuns["user_id"] = data_shixuns["user_id"].astype("float").astype("int")
map_list = {
'Java': 'Java',
'MySQL/Java': 'MySQL',
'Python': 'Python',
'C++': 'C/C++',
'Python3.6': 'Python',
'JFinal': 'JFinal',
'Python2.7': 'Python',
'Dynamips': 'Dynamips',
'Ethereum': 'Ethereum',
'Html': 'Html',
'MachineLearning': 'MachineLearning',
'Docker': 'Docker',
'C': 'C/C++',
'MySQL/Python3.6': 'MySQL',
'Verilog': 'Verilog',
'PHP/Web': 'PHP/Web',
'Android': 'Android',
'Golang': 'Golang',
'Hadoop': 'Hadoop',
'Matlab': 'Matlab',
'Shell': 'Shell',
'Git': 'Git',
'Ruby': 'Ruby',
'Perl6': 'Perl6',
'Kotlin': 'Kotlin',
'JavaScript': 'JavaScript',
'Other': 'Other'
}
data_shixuns['language'] = data_shixuns['language'].map(map_list)
return data_shixuns
def get_program_data():
"""
生成用户编程训练数据
"""
data = read_action_data()
data = data[['user_id', 'language']]
# 计算用户在每一类语言实训的总数
data['value'] = 1
data = data.pivot_table(index='user_id', columns='language', values='value', aggfunc='count').fillna(0)
data = data.astype('int')
# 计算用户每类语言的使用占比
data = data.div(data.sum(axis=0), axis=1)
return data
if __name__ == '__main__':
# 取数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1,22 @@
import pickle
from config import logger
from config import test_user_id
from config import programming_ability_analysis_path
logger.info('加载用户用户编程能力字典')
user_programming_ability_dict = pickle.load(open(programming_ability_analysis_path + 'results/user_programming_ability_dict.pkl', 'rb'))
def user_programming_ability_predict(user_id):
"""
用户编程能力预测
"""
if user_id not in user_programming_ability_dict:
result = -1
else:
result = user_programming_ability_dict[user_id]
return result
if __name__ == '__main__':
result = user_programming_ability_predict(user_id=test_user_id)
print('用户ID', test_user_id, '用户编程能力:', result)

View File

@ -0,0 +1,2 @@
## 这里面存放模型训练结果

View File

@ -0,0 +1,82 @@
import os
import numpy as np
import pandas as pd
import pickle
from sklearn.decomposition import PCA
from sklearn.cluster import KMeans
from config import logger
from config import RANDOM_SEED
from data_process import get_program_data
from config import programming_ability_analysis_path
def iflabel(x):
if x == "高高高":
return 5
elif x == "高低高":
return 4
elif x == "高高低":
return 4
elif x == "高低低":
return 3
elif x == "低高高":
return 4
elif x == "低低高":
return 2
elif x == "低高低":
return 2
elif x == "低低低":
return 1
def train():
"""
用户编程数据分析模型训练
"""
logger.info('开始训练用编程能力分析模型...')
data = get_program_data()
# 跟据PCA算法计算出影响力最大的三种语言
pca_model = PCA(n_components=3)
pca_data = pca_model.fit_transform(data)
new_data = pd.DataFrame(np.array(pca_data), index=range(len(list(pca_data))), columns=['1', '2', '3'])
new_data['user_id'] = list(new_data.index)
# 根据这三种语言给用户进行聚类
df_cluster = new_data[['user_id', '1', '2', '3']]
kmeans_model = KMeans(n_clusters=8, random_state=RANDOM_SEED).fit(df_cluster[['1', '2', '3']])
df_cluster["label"] = kmeans_model.labels_
# 计算每一类的均值,超过均值设置为'高',否则设置为'低'
df_mean = df_cluster[['1', '2', '3', 'label']].groupby("label").agg("mean")
df_mean["IF1"] = df_mean["1"].apply(lambda x: "" if x > df_mean["1"].mean() else "")
df_mean["IF2"] = df_mean["2"].apply(lambda x: "" if x > df_mean["2"].mean() else "")
df_mean["IF3"] = df_mean["3"].apply(lambda x: "" if x > df_mean["3"].mean() else "")
# 根据打分函数给用户进行打分
df_mean["temp"] = df_mean["IF1"] + df_mean["IF2"] + df_mean["IF3"]
df_mean["label"] = df_mean["temp"].apply(lambda x: iflabel(x))
df_cluster["IF1"] = df_cluster["1"].apply(lambda x: "" if x > df_mean["1"].mean() else "")
df_cluster["IF2"] = df_cluster["2"].apply(lambda x: "" if x > df_mean["1"].mean() else "")
df_cluster["IF3"] = df_cluster["3"].apply(lambda x: "" if x > df_mean["1"].mean() else "")
df_cluster["temp"] = df_cluster["IF1"] + df_cluster["IF2"] + df_cluster["IF3"]
df_cluster["label"] = df_cluster["temp"].apply(lambda x: iflabel(x))
# 将用户 id 和编程能力得分提取存储起来
df_cluster = df_cluster.rename(columns={'label': 'programming_ability'})
result_csv = df_cluster[['user_id','programming_ability']]
result_csv['user_id'] = result_csv['user_id'].astype(int)
result_csv['programming_ability'] = result_csv['programming_ability'].astype(int)
results_dict = dict(zip(result_csv['user_id'], result_csv['programming_ability']))
# 保存模型和结果
pickle.dump(pca_model, open(programming_ability_analysis_path + 'results/user_programming_ability_pca_model.pkl', 'wb'))
pickle.dump(kmeans_model, open(programming_ability_analysis_path + 'results/user_programming_ability_cluster_model.pkl', 'wb'))
pickle.dump(results_dict, open(programming_ability_analysis_path + 'results/user_programming_ability_dict.pkl', 'wb'))
logger.info('用户编程能力分析模型训练完成')
return results_dict
if __name__ == '__main__':
train()

View File

@ -0,0 +1,103 @@
import argparse
import os
from config import logger
def del_output_file(parent_path):
"""
删除之前的数据和输出
"""
files_list = os.listdir(parent_path)
for file_name in files_list:
if not os.path.isdir(file_name) and ('.py' not in file_name):
file_name = parent_path + file_name
if os.path.exists(file_name):
logger.info('删除文件: ' + file_name)
os.system('rm ' + file_name)
if __name__ == '__main__':
parser = argparse.ArgumentParser()
parser.add_argument('--silent', action='store_true', help='是否以silent模式运行')
args = parser.parse_args()
silent_mode = args.silent
if silent_mode:
confirm = 'y'
else:
confirm = input('确认要重新开始训练所有模型吗y/n' + '\n').strip()
if confirm == 'y' or confirm == 'Y':
if silent_mode:
confirm = 'y'
else:
confirm = input('是否删除之前的数据和所有输出y/n' + '\n').strip()
if confirm == 'y' or confirm == 'Y':
parent_path_list = []
parent_path_list.clear()
parent_path_list.append('./activity_analysis/data/')
parent_path_list.append('./activity_analysis/results/')
parent_path_list.append('./contribution_analysis/data/')
parent_path_list.append('./contribution_analysis/results/')
parent_path_list.append('./interests_analysis/data/')
parent_path_list.append('./interests_analysis/results/')
parent_path_list.append('./learning_ability_analysis/data/')
parent_path_list.append('./learning_ability_analysis/results/')
parent_path_list.append('./professional_ability_analysis/data/')
parent_path_list.append('./professional_ability_analysis/results/')
parent_path_list.append('./programming_ability_analysis/data/')
parent_path_list.append('./programming_ability_analysis/results/')
parent_path_list.append('./user_label_analysis/data/')
parent_path_list.append('./user_label_analysis/results/')
for parent_path in parent_path_list:
if not os.path.exists(parent_path):
os.mkdir(parent_path)
del_output_file(parent_path)
# 用户活跃度分析数据处理
os.system('python ./activity_analysis/data_process.py')
# 用户活跃度分析模型训练
os.system('python ./activity_analysis/train.py')
# 用户贡献度分析数据处理
os.system('python ./contribution_analysis/data_process.py')
# 用户贡献度分析模型训练
os.system('python ./contribution_analysis/train.py')
# 用户爱好分析数据处理
os.system('python ./interests_analysis/data_process.py')
# 用户爱好分析模型训练
os.system('python ./interests_analysis/train.py')
# 用户学习能力分析数据处理
os.system('python ./learning_ability_analysis/data_process.py')
# 用户学习能力分析模型训练
os.system('python ./learning_ability_analysis/train.py')
# 用户专业能力分析数据处理
os.system('python ./professional_ability_analysis/data_process.py')
# 用户专业能力分析模型训练
os.system('python ./professional_ability_analysis/train.py')
# 用户编程能力分析数据处理
os.system('python ./programming_ability_analysis/data_process.py')
# 用户编程能力分析模型训练
os.system('python ./programming_ability_analysis/train.py')
# 用户标签分析数据处理
os.system('python ./user_label_analysis/data_process.py')
# 用户标签分析模型训练
os.system('python ./user_label_analysis/train.py')
# 启动模型预测接口
os.system('python ./flask_app.py')

View File

@ -0,0 +1,2 @@
## 这里面存放的是数据

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

View File

@ -0,0 +1,153 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
import pickle
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from utils import get_before_date
from config import data_before_days
from config import activity_analysis_path
from config import contribution_analysis_path
from config import user_label_analysis_path
from config import learning_ability_analysis_path
from config import professional_ability_analysis_path
from config import programming_ability_analysis_path
maplist = {
"上海":"华东",
"江苏":"华东",
"浙江":"华东",
"安徽":"华东",
"福建":"华东",
"江西":"华东",
"山东":"华东",
"台湾":"华东",
"北京":"华北",
"天津":"华北",
"河北":"华北",
"山西":"华北",
"内蒙古":"华北",
"广东":"华南",
"广西":"华南",
"海南":"华南",
"香港":"华南",
"澳门":"华南",
"河南":"华中",
"湖北":"华中",
"湖南":"华中",
"陕西":"西北",
"甘肃":"西北",
"青海":"西北",
"宁夏":"西北",
"新疆":"西北"
}
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户标签数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text1 = f"""
SELECT id, name, province
FROM schools
"""
data_schools = pd.read_sql(sql_text1, con=sql_conn)
data_schools.to_csv( user_label_analysis_path + 'data/schools.csv', index=False, header=True, sep='\t')
logger.info("用户标签数据1下载完毕" + str(data_schools.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text2 = f"""
SELECT user_id,username,school_name
FROM subject_user_infos
"""
data_user_infos = pd.read_sql(sql_text2, con=sql_conn)
data_user_infos.to_csv(user_label_analysis_path + 'data/subject_user_infos.csv', index=False, header=True, sep='\t')
logger.info(
"用户标签数据2下载完毕" + str(data_user_infos.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理用户标签数据
"""
# 提取学校数据并将学校所在地区的分区标上
data_school = pd.read_csv(user_label_analysis_path + 'data/schools.csv', sep='\t')
data_school = data_school[['id', 'name', 'province']]
data_school = data_school.dropna()
data_school['area'] = data_school['province'].map(maplist)
data_user_infos = pd.read_csv(user_label_analysis_path + 'data/subject_user_infos.csv', sep='\t')
data_user_infos = data_user_infos[['user_id', 'username', 'school_name']]
data_user_infos = data_user_infos.drop_duplicates(subset='user_id')
return data_school, data_user_infos
def get_label_data():
"""
生成用户排序数据
"""
data_school, data_user_infos = read_action_data()
# 提取之前所有用户各项指标的得分
user_activity_dict = pickle.load(open(activity_analysis_path + 'results/user_activity_dict.pkl', 'rb'))
user_contribution_dict = pickle.load(open(contribution_analysis_path + 'results/user_contribution_dict.pkl', 'rb'))
user_learning_ability_dict = pickle.load(
open(learning_ability_analysis_path + 'results/user_learning_ability_dict.pkl', 'rb'))
user_professional_ability_dict = pickle.load(
open(professional_ability_analysis_path + 'results/user_professional_ability_dict.pkl', 'rb'))
user_programming_ability_dict = pickle.load(
open(programming_ability_analysis_path + 'results/user_programming_ability_dict.pkl', 'rb'))
df_user_activity = pd.DataFrame.from_dict(user_activity_dict, orient='index', columns=['activity'])
df_user_activity = df_user_activity.reset_index()
df_user_activity = df_user_activity.rename(columns={'index': 'user_id'})
df_user_contribution = pd.DataFrame.from_dict(user_contribution_dict, orient='index', columns=['contribution'])
df_user_contribution = df_user_contribution.reset_index()
df_user_contribution = df_user_contribution.rename(columns={'index': 'user_id'})
df_learning_ability = pd.DataFrame.from_dict(user_learning_ability_dict, orient='index',
columns=['learning_ability'])
df_learning_ability = df_learning_ability.reset_index()
df_learning_ability = df_learning_ability.rename(columns={'index': 'user_id'})
df_professional_ability = pd.DataFrame.from_dict(user_professional_ability_dict, orient='index',
columns=['professional_ability'])
df_professional_ability = df_professional_ability.reset_index()
df_professional_ability = df_professional_ability.rename(columns={'index': 'user_id'})
df_programming_ability = pd.DataFrame.from_dict(user_programming_ability_dict, orient='index',
columns=['programming_ability'])
df_programming_ability = df_programming_ability.reset_index()
df_programming_ability = df_programming_ability.rename(columns={'index': 'user_id'})
# 将未统计到用户得分的项目都给他打1分
df_merge = pd.merge(df_user_activity, df_user_contribution, on='user_id', how='outer')
df_merge = df_merge.fillna(1)
df_merge = pd.merge(df_merge, df_learning_ability, on='user_id', how='outer')
df_merge = df_merge.fillna(1)
df_merge = pd.merge(df_merge, df_professional_ability, on='user_id', how='outer')
df_merge = df_merge.fillna(1)
df_merge = pd.merge(df_merge, df_programming_ability, on='user_id', how='outer')
df_merge = df_merge.fillna(1)
return df_merge, data_school, data_user_infos
if __name__ == '__main__':
# 取数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1,170 @@
import os
import numpy as np
import pandas as pd
from config import logger
import datetime
import pymysql
import pickle
from config import mysql_database, mysql_passwd
from config import mysql_host, mysql_port, mysql_user
from utils import get_before_date
from config import data_before_days
from config import activity_analysis_path
from config import contribution_analysis_path
from config import user_label_analysis_path
from config import learning_ability_analysis_path
from config import professional_ability_analysis_path
from config import programming_ability_analysis_path
maplist = {
"上海":"华东",
"江苏":"华东",
"浙江":"华东",
"安徽":"华东",
"福建":"华东",
"江西":"华东",
"山东":"华东",
"台湾":"华东",
"北京":"华北",
"天津":"华北",
"河北":"华北",
"山西":"华北",
"内蒙古":"华北",
"广东":"华南",
"广西":"华南",
"海南":"华南",
"香港":"华南",
"澳门":"华南",
"河南":"华中",
"湖北":"华中",
"湖南":"华中",
"陕西":"西北",
"甘肃":"西北",
"青海":"西北",
"宁夏":"西北",
"新疆":"西北"
}
def get_data_from_mysql(latest_data_date):
"""
从mysql数据库获取原始数据
"""
start = datetime.datetime.now()
logger.info("开始获取用户标签数据...")
sql_conn = pymysql.connect(host=mysql_host,
user=mysql_user,
passwd=mysql_passwd,
port=mysql_port,
db=mysql_database)
sql_text1 = f"""
select
t1.id,
t1.name,
t1.province
from
schools t1
left join user_extensions t2 on t1.id = t2.school_id
left join users t3 on t2.user_id = t3.id
where
DATE_FORMAT(t3.last_login_on, '%Y-%m-%d') >= '{latest_data_date}'
and t1.province is not null
and t1.province != ''
"""
data_schools = pd.read_sql(sql_text1, con=sql_conn)
data_schools.to_csv( user_label_analysis_path + 'data/schools.csv', index=False, header=True, sep='\t')
logger.info("用户标签数据1下载完毕" + str(data_schools.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
sql_text2 = f"""
SELECT
user_id,
username,
school_name
FROM
subject_user_infos t1
LEFT JOIN users t2 ON t1.user_id = t2.id
WHERE
DATE_FORMAT(t2.last_login_on, '%Y-%m-%d') >= '{latest_data_date}'
"""
data_user_infos = pd.read_sql(sql_text2, con=sql_conn)
data_user_infos.to_csv(user_label_analysis_path + 'data/subject_user_infos.csv', index=False, header=True, sep='\t')
logger.info(
"用户标签数据2下载完毕" + str(data_user_infos.shape[0]) + "条数据,总耗时" + str((datetime.datetime.now() - start).seconds) + "")
def read_action_data():
"""
处理用户标签数据
"""
# 提取学校数据并将学校所在地区的分区标上
data_school = pd.read_csv(user_label_analysis_path + 'data/schools.csv', sep='\t')
data_school = data_school[['id', 'name', 'province']]
data_school = data_school.dropna()
data_school['area'] = data_school['province'].map(maplist)
data_user_infos = pd.read_csv(user_label_analysis_path + 'data/subject_user_infos.csv', sep='\t')
data_user_infos = data_user_infos[['user_id', 'username', 'school_name']]
data_user_infos = data_user_infos.drop_duplicates(subset='user_id')
return data_school, data_user_infos
def get_label_data():
"""
生成用户排序数据
"""
data_school, data_user_infos = read_action_data()
# 提取之前所有用户各项指标的得分
user_activity_dict = pickle.load(open(activity_analysis_path + 'results/user_activity_dict.pkl', 'rb'))
user_contribution_dict = pickle.load(open(contribution_analysis_path + 'results/user_contribution_dict.pkl', 'rb'))
user_learning_ability_dict = pickle.load(
open(learning_ability_analysis_path + 'results/user_learning_ability_dict.pkl', 'rb'))
user_professional_ability_dict = pickle.load(
open(professional_ability_analysis_path + 'results/user_professional_ability_dict.pkl', 'rb'))
user_programming_ability_dict = pickle.load(
open(programming_ability_analysis_path + 'results/user_programming_ability_dict.pkl', 'rb'))
df_user_activity = pd.DataFrame.from_dict(user_activity_dict, orient='index', columns=['activity'])
df_user_activity = df_user_activity.reset_index()
df_user_activity = df_user_activity.rename(columns={'index': 'user_id'})
df_user_contribution = pd.DataFrame.from_dict(user_contribution_dict, orient='index', columns=['contribution'])
df_user_contribution = df_user_contribution.reset_index()
df_user_contribution = df_user_contribution.rename(columns={'index': 'user_id'})
df_learning_ability = pd.DataFrame.from_dict(user_learning_ability_dict, orient='index',
columns=['learning_ability'])
df_learning_ability = df_learning_ability.reset_index()
df_learning_ability = df_learning_ability.rename(columns={'index': 'user_id'})
df_professional_ability = pd.DataFrame.from_dict(user_professional_ability_dict, orient='index',
columns=['professional_ability'])
df_professional_ability = df_professional_ability.reset_index()
df_professional_ability = df_professional_ability.rename(columns={'index': 'user_id'})
df_programming_ability = pd.DataFrame.from_dict(user_programming_ability_dict, orient='index',
columns=['programming_ability'])
df_programming_ability = df_programming_ability.reset_index()
df_programming_ability = df_programming_ability.rename(columns={'index': 'user_id'})
# 将未统计到用户得分的项目都给他打1分
df_merge = pd.merge(df_user_activity, df_user_contribution, on='user_id', how='outer')
df_merge = df_merge.fillna(1)
df_merge = pd.merge(df_merge, df_learning_ability, on='user_id', how='outer')
df_merge = df_merge.fillna(1)
df_merge = pd.merge(df_merge, df_professional_ability, on='user_id', how='outer')
df_merge = df_merge.fillna(1)
df_merge = pd.merge(df_merge, df_programming_ability, on='user_id', how='outer')
df_merge = df_merge.fillna(1)
return df_merge, data_school, data_user_infos
if __name__ == '__main__':
# 取数据
latest_data_date = get_before_date(data_before_days)
get_data_from_mysql(latest_data_date)

View File

@ -0,0 +1 @@
## 这里是模型的代码说明

View File

@ -0,0 +1,36 @@
import pickle
from config import logger
from config import test_user_id
from config import user_label_analysis_path
logger.info('加载用户排序字典')
results_dict_rank = pickle.load(open(user_label_analysis_path + 'results/user_rank_dict.pkl', 'rb'))
results_dict_area_rank = pickle.load(open(user_label_analysis_path + 'results/user_area_rank_dict.pkl', 'rb'))
results_dict_school_rank = pickle.load(open(user_label_analysis_path+ 'results/user_school_rank_dict.pkl', 'rb'))
results_dict_label = pickle.load(open(user_label_analysis_path + 'results/user_label_dict.pkl', 'rb'))
def user_label_predict(user_id):
"""
用户贡献度预测
"""
if user_id not in results_dict_rank:
results_rank = -1
results_area_rank = -1
results_school_rank = -1
results_label = -1
else:
results_rank = results_dict_rank[user_id]
results_area_rank = results_dict_area_rank[user_id]
results_school_rank = results_dict_school_rank[user_id]
results_label = results_dict_label[user_id]
return results_rank, results_area_rank, results_school_rank, results_label
if __name__ == '__main__':
results_rank, results_area_rank, results_school_rank, results_label = user_label_predict(user_id=test_user_id)
print('用户ID', test_user_id,
'全国排名:', results_rank,
'地区排名:',results_area_rank,
'学校排名:',results_school_rank,
'用户标签:',results_label)

View File

@ -0,0 +1,2 @@
## 这里面存放模型训练结果

View File

@ -0,0 +1,128 @@
import os
import numpy as np
import pandas as pd
import pickle
from sklearn.cluster import KMeans
from config import logger
from config import RANDOM_SEED
from data_process import get_label_data
from config import user_label_analysis_path
from learning2rank.rank import RankNet
import chainer
# 根据比例打标签
def iflabel(x,data):
class_names = ['学霸型', '学习优秀型','学习进步型' ,'学习懒散型']
num = len(data)
if x >= 0 and x < num*0.15:
return class_names[0]
elif x >= num*0.15 and x < num*0.5:
return class_names[1]
elif x >= num*0.5 and x < num*0.85:
return class_names[2]
else:
return class_names[3]
def train():
"""
用户排序模型训练
"""
# 根据 5 想指标计算总和并根据总和将所有用户切分为 4 类作为用户标签
logger.info('开始训练用户排序分析模型...')
data, data_shool, data_user = get_label_data()
data['label'] = data[
['activity', 'contribution', 'learning_ability', 'professional_ability', 'programming_ability']].sum(1)
q = pd.cut(data.label, 4, labels=[1, 2, 3, 4])
data['label'] = q
'''
此处为另一种思路的代码采用k_means聚类
'''
# cluster = data[['activity','contribution','learning_ability','professional_ability','programming_ability']]
# kmeans = KMeans(n_clusters=5, random_state=0).fit(cluster)
# data["label"] = kmeans.labels_ + 1
# 将用户 5 个维度的数据作为特征,上面切分的类别作为标签送入排序模型进行训练然后到处预测结果
rank_model = RankNet.RankNet()
X = np.array(data.iloc[:, 1:6].values)
y = np.array(data.iloc[:, -1].values)
rank_model.fit(X, y)
user_rank = rank_model.predict(X)
# 根据预测结果将分值转化为排名
ranks_ = []
for i in range(len(user_rank)):
ranks_.append(user_rank[i][0])
ranks1 = pd.Series(ranks_, data.index)
data['RankNet_predict_value'] = ranks1
# 根据排名将用户设置为4种类型'学习懒散型', '学习进步型', '学习优秀型', '学霸型'
data = data.sort_values(by='RankNet_predict_value', ascending=False)
data['rank'] = data['RankNet_predict_value'].rank(ascending=False, method='min')
data['rank'] = data['rank'].astype('float').astype('int')
data = data.reset_index(drop=True)
data['index'] = data.index
data['final_label'] = data['index'].apply(lambda x: iflabel(x, data))
data.drop('index', axis=1, inplace=True)
# 提取用户 id ,姓名,学校,地区
data_user_info_1 = pd.merge(data, data_user, on='user_id')
data_user_info_2 = data_shool[['name', 'area']]
data_user_info_2 = data_user_info_2.rename(columns={'name': 'school_name'})
data_user_infos = pd.merge(data_user_info_1, data_user_info_2, on='school_name')
# 根据所属地区,将每个地区的用户的排名
areas = list(data_user_infos.area.unique())
df_area_rank = data_user_infos.copy()
df_area_rank['area_rank'] = 0
df_area_rank = df_area_rank.drop(df_area_rank.index[0:len(df_area_rank)], 0)
for area in areas:
df_user_copy = data_user_infos[data_user_infos.area == area].copy()
df_user_copy['area_rank'] = df_user_copy['RankNet_predict_value'].rank(ascending=False, method='min')
df_user_copy['area_rank'] = df_user_copy['area_rank'].astype('int')
df_area_rank = pd.concat([df_user_copy, df_area_rank])
# 根据学校名称,将每个学校的用户进行内部排名
schools = list(data_user_infos.school_name.unique())
df_school_rank = df_area_rank.copy()
df_school_rank['school_rank'] = 0
df_school_rank = df_school_rank.drop(df_school_rank.index[0:len(df_school_rank)], 0)
for school in schools:
df_school_copy = df_area_rank[df_area_rank.school_name == school].copy()
df_school_copy['school_rank'] = df_school_copy['RankNet_predict_value'].rank(ascending=False, method='min')
df_school_copy['school_rank'] = df_school_copy['school_rank'].astype('int')
df_school_rank = pd.concat([df_school_copy, df_school_rank])
# 将用户 id ,全国排名,标签,地区排名和学校排名提取并且保存
result_csv = df_school_rank[['user_id',
'rank',
'final_label',
'area_rank',
'school_rank']]
df_school_rank.to_csv(user_label_analysis_path + 'results/user_label_analysis.csv',columns=["user_id","activity"
,"contribution","learning_ability","professional_ability","programming_ability","label","RankNet_predict_value","rank",
"final_label","username","school_name"])
result_csv['user_id'] = result_csv['user_id'].astype(int)
result_csv['rank'] = result_csv['rank'].astype(int)
result_csv['area_rank'] = result_csv['area_rank'].astype(int)
result_csv['school_rank'] = result_csv['school_rank'].astype(int)
results_dict_rank = dict(zip(result_csv['user_id'], result_csv['rank']))
results_dict_area_rank = dict(zip(result_csv['user_id'], result_csv['area_rank']))
results_dict_school_rank = dict(zip(result_csv['user_id'], result_csv['school_rank']))
results_dict_label = dict(zip(result_csv['user_id'], result_csv['final_label']))
# # 保存模型和结果
pickle.dump(rank_model, open(user_label_analysis_path + 'results/user_label_rankmodel_dict.pkl', 'wb'))
pickle.dump(results_dict_rank, open(user_label_analysis_path + 'results/user_rank_dict.pkl', 'wb'))
pickle.dump(results_dict_area_rank, open(user_label_analysis_path + 'results/user_area_rank_dict.pkl', 'wb'))
pickle.dump(results_dict_school_rank, open(user_label_analysis_path + 'results/user_school_rank_dict.pkl', 'wb'))
pickle.dump(results_dict_label, open(user_label_analysis_path + 'results/user_label_dict.pkl', 'wb'))
logger.info('用户排序模型训练完成')
return results_dict_rank, results_dict_area_rank, results_dict_school_rank, results_dict_label
if __name__ == '__main__':
train()

View File

@ -0,0 +1,66 @@
import os
import logging
import random
import torch
import numpy as np
import datetime
def create_logger(log_path):
"""
将日志输出到日志文件和控制台
"""
logger = logging.getLogger(__name__)
logger.setLevel(logging.INFO)
formatter = logging.Formatter('%(asctime)s - %(levelname)s - %(message)s')
# 创建一个handler用于写入日志文件
file_handler = logging.FileHandler(filename=log_path)
file_handler.setFormatter(formatter)
file_handler.setLevel(logging.INFO)
if file_handler not in logger.handlers:
logger.addHandler(file_handler)
# 创建一个handler用于将日志输出到控制台
console = logging.StreamHandler()
console.setLevel(logging.DEBUG)
console.setFormatter(formatter)
if console not in logger.handlers:
logger.addHandler(console)
return logger
def get_file_name(fname):
"""
获取文件名
"""
return os.path.split(fname)[-1].split(".")[0]
def get_file_size(fname):
"""
获取文件大小MB
"""
fsize = os.path.getsize(fname)
fsize = fsize/float(1024 * 1024)
return round(fsize, 2)
def set_seed(seed):
"""
设置随机数种子
"""
random.seed(seed)
np.random.seed(seed)
torch.manual_seed(seed)
if torch.cuda.is_available():
torch.cuda.manual_seed_all(seed)
def get_before_date(n):
"""
获取前N天的日期
"""
today = datetime.datetime.now()
# 计算偏移量
offset = datetime.timedelta(days=-n)
re_date = (today + offset).strftime('%Y-%m-%d')
return re_date

74
utils.py Normal file
View File

@ -0,0 +1,74 @@
import os
import logging
import datetime
import psycopg2
import calendar
def create_logger(log_path):
"""
将日志输出到日志文件和控制台
"""
logger = logging.getLogger(__name__)
logger.setLevel(logging.INFO)
formatter = logging.Formatter('%(asctime)s - %(levelname)s - %(message)s')
# 创建一个handler用于写入日志文件
file_handler = logging.FileHandler(filename=log_path)
file_handler.setFormatter(formatter)
file_handler.setLevel(logging.INFO)
if file_handler not in logger.handlers:
logger.addHandler(file_handler)
# 创建一个handler用于将日志输出到控制台
console = logging.StreamHandler()
console.setLevel(logging.DEBUG)
console.setFormatter(formatter)
if console not in logger.handlers:
logger.addHandler(console)
return logger
def get_conn():
from config import host, port, database, user, passwd
conn_string = "host=" + host + " port=" + str(port) + " dbname=" + database + " user=" + user + " password=" + passwd
conn = psycopg2.connect(conn_string)
return conn
def get_month_date(year, month):
date_list = []
for i in range(calendar.monthrange(year, month)[1] + 1)[1:]:
str1 = str(year) + '-' + str("%02d" % month) + '-' + str("%02d" % i)
date_list.append(str1)
return date_list
def get_file_name(fname):
"""
获取文件名
"""
return os.path.split(fname)[-1].split(".")[0]
def get_file_size(fname):
"""
获取文件大小MB
"""
fsize = os.path.getsize(fname)
fsize = fsize/float(1024 * 1024)
return round(fsize, 2)
def get_before_date(n):
"""
获取前N天的日期
"""
today = datetime.datetime.now()
# 计算偏移量
offset = datetime.timedelta(days=-n)
re_date = (today + offset).strftime('%Y-%m-%d')
return re_date