Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
A
Amazon-Selection-Data
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
abel_cjy
Amazon-Selection-Data
Commits
2a408ea3
Commit
2a408ea3
authored
Aug 19, 2026
by
chenyuanjie
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
流量选品-品牌判断内部asin
parent
195c93a8
Show whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
34 additions
and
8 deletions
+34
-8
dim_asin_detail.py
Pyspark_job/dim/dim_asin_detail.py
+19
-3
kafka_flow_asin_detail_to_doris.py
Pyspark_job/my_kafka/kafka_flow_asin_detail_to_doris.py
+15
-5
No files found.
Pyspark_job/dim/dim_asin_detail.py
View file @
2a408ea3
"""
"""
@Author :
HuangJian
@Author :
CT
@Description : 关键词与Asin详情维表
@Description : 关键词与Asin详情维表
@SourceTable :
@SourceTable :
①ods_asin_keep_date
①ods_asin_keep_date
...
@@ -66,6 +66,7 @@ class DimAsinDetail(object):
...
@@ -66,6 +66,7 @@ class DimAsinDetail(object):
self
.
df_asin_new_cate
=
self
.
spark
.
sql
(
"select 1+1;"
)
self
.
df_asin_new_cate
=
self
.
spark
.
sql
(
"select 1+1;"
)
self
.
df_user_package_num
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
self
.
df_user_package_num
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
self
.
df_self_asin
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
self
.
df_self_asin
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
self
.
df_self_brand
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
self
.
df_asin_category
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
self
.
df_asin_category
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
self
.
df_asin_variat
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
self
.
df_asin_variat
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
self
.
df_keepa_tracking
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
self
.
df_keepa_tracking
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
...
@@ -195,6 +196,17 @@ class DimAsinDetail(object):
...
@@ -195,6 +196,17 @@ class DimAsinDetail(object):
df_self_asin
=
self
.
spark
.
sql
(
sqlQuery
=
sql
)
df_self_asin
=
self
.
spark
.
sql
(
sqlQuery
=
sql
)
self
.
df_self_asin
=
F
.
broadcast
(
df_self_asin
)
self
.
df_self_asin
=
F
.
broadcast
(
df_self_asin
)
self
.
df_self_asin
.
show
(
10
,
truncate
=
False
)
self
.
df_self_asin
.
show
(
10
,
truncate
=
False
)
print
(
"5b.读取amazon_brand,获得内部品牌信息(品牌级内部asin判断)"
)
mysql_brand_sql
=
"""
SELECT DISTINCT LOWER(TRIM(brand_name)) AS asin_brand_name FROM amazon_brand WHERE brand_type = '1'
"""
mysql_con
=
DBUtil
.
get_connection_info
(
"mysql"
,
"us"
)
df_self_brand
=
SparkUtil
.
read_jdbc_query
(
session
=
self
.
spark
,
url
=
mysql_con
[
'url'
],
pwd
=
mysql_con
[
'pwd'
],
username
=
mysql_con
[
'username'
],
query
=
mysql_brand_sql
)
.
withColumn
(
"asin_is_self_brand"
,
F
.
lit
(
1
))
self
.
df_self_brand
=
F
.
broadcast
(
df_self_brand
)
self
.
df_self_brand
.
show
(
10
,
truncate
=
False
)
print
(
"6. node_id对应的头部分类信息"
)
print
(
"6. node_id对应的头部分类信息"
)
self
.
df_asin_new_cate
=
get_node_first_id_df
(
self
.
site_name
,
self
.
spark
)
self
.
df_asin_new_cate
=
get_node_first_id_df
(
self
.
site_name
,
self
.
spark
)
self
.
df_asin_new_cate
=
self
.
df_asin_new_cate
.
filter
(
'node_id is not null'
)
.
persist
(
StorageLevel
.
DISK_ONLY
)
self
.
df_asin_new_cate
=
self
.
df_asin_new_cate
.
filter
(
'node_id is not null'
)
.
persist
(
StorageLevel
.
DISK_ONLY
)
...
@@ -452,10 +464,14 @@ class DimAsinDetail(object):
...
@@ -452,10 +464,14 @@ class DimAsinDetail(object):
df_alarm_brand
=
df_alarm_brand
.
repartition
(
100
)
df_alarm_brand
=
df_alarm_brand
.
repartition
(
100
)
self
.
df_asin_detail
=
self
.
df_asin_detail
.
join
(
df_alarm_brand
,
on
=
[
'asin_brand_name'
],
how
=
'left'
)
self
.
df_asin_detail
=
self
.
df_asin_detail
.
join
(
df_alarm_brand
,
on
=
[
'asin_brand_name'
],
how
=
'left'
)
self
.
df_asin_detail
=
self
.
df_asin_detail
.
na
.
fill
({
"asin_is_alarm"
:
0
})
self
.
df_asin_detail
=
self
.
df_asin_detail
.
na
.
fill
({
"asin_is_alarm"
:
0
})
# 处理是否内部asin信息
# 处理是否内部asin信息
:Hive ods_self_asin(asin级) 或 amazon_brand(品牌级,brand_type=1) 命中任一即为内部asin
self
.
df_asin_detail
=
self
.
df_asin_detail
.
join
(
self
.
df_self_asin
,
on
=
[
'asin'
],
how
=
'left'
)
self
.
df_asin_detail
=
self
.
df_asin_detail
.
join
(
self
.
df_self_asin
,
on
=
[
'asin'
],
how
=
'left'
)
self
.
df_asin_detail
=
self
.
df_asin_detail
.
na
.
fill
({
"asin_is_self"
:
0
})
self
.
df_asin_detail
=
self
.
df_asin_detail
.
join
(
self
.
df_self_brand
,
on
=
[
'asin_brand_name'
],
how
=
'left'
)
self
.
df_asin_detail
=
self
.
df_asin_detail
.
withColumn
(
"asin_is_self"
,
F
.
when
((
F
.
col
(
"asin_is_self"
)
==
1
)
|
(
F
.
col
(
"asin_is_self_brand"
)
==
1
),
F
.
lit
(
1
))
.
otherwise
(
F
.
lit
(
0
))
)
.
drop
(
"asin_is_self_brand"
)
self
.
df_self_asin
.
unpersist
()
self
.
df_self_asin
.
unpersist
()
self
.
df_self_brand
.
unpersist
()
# 处理影视标签字段
# 处理影视标签字段
def
handle_asin_label
(
self
):
def
handle_asin_label
(
self
):
...
...
Pyspark_job/my_kafka/kafka_flow_asin_detail_to_doris.py
View file @
2a408ea3
...
@@ -113,6 +113,7 @@ class KafkaFlowAsinDetail(Templates):
...
@@ -113,6 +113,7 @@ class KafkaFlowAsinDetail(Templates):
self
.
df_asin_new_cate
=
self
.
spark
.
sql
(
"select 1+ 1;"
)
self
.
df_asin_new_cate
=
self
.
spark
.
sql
(
"select 1+ 1;"
)
self
.
df_asin_category
=
self
.
spark
.
sql
(
"select 1+1;"
)
self
.
df_asin_category
=
self
.
spark
.
sql
(
"select 1+1;"
)
self
.
df_self_asin
=
self
.
spark
.
sql
(
"select 1+1;"
)
self
.
df_self_asin
=
self
.
spark
.
sql
(
"select 1+1;"
)
self
.
df_self_brand
=
self
.
spark
.
sql
(
"select 1+1;"
)
self
.
df_hide_category
=
self
.
spark
.
sql
(
"select 1+1;"
)
self
.
df_hide_category
=
self
.
spark
.
sql
(
"select 1+1;"
)
self
.
color_set
=
set
()
self
.
color_set
=
set
()
# udf函数注册
# udf函数注册
...
@@ -634,7 +635,9 @@ class KafkaFlowAsinDetail(Templates):
...
@@ -634,7 +635,9 @@ class KafkaFlowAsinDetail(Templates):
)
)
df
=
df
.
drop
(
"asin_label_list"
,
"badge_list"
)
df
=
df
.
drop
(
"asin_label_list"
,
"badge_list"
)
# 4. asin_type(对齐 dwt.handle_asin_is_hide:内部asin=1 / 数字虚拟类目=2 / 隐藏分类=3 / 默认0)
# 4. asin_type(对齐 dwt.handle_asin_is_hide:内部asin=1 / 数字虚拟类目=2 / 隐藏分类=3 / 默认0)
# 内部asin判断两步:asin级(ods_self_asin) 或 品牌级(amazon_brand,brand_type=1) 命中任一即可
df
=
df
.
join
(
self
.
df_self_asin
,
on
=
[
'asin'
],
how
=
'left'
)
df
=
df
.
join
(
self
.
df_self_asin
,
on
=
[
'asin'
],
how
=
'left'
)
df
=
df
.
join
(
self
.
df_self_brand
,
on
=
[
'brand'
],
how
=
'left'
)
df
=
df
.
join
(
self
.
df_hide_category
,
on
=
[
'asin_bs_cate_current_id'
],
how
=
'left'
)
df
=
df
.
join
(
self
.
df_hide_category
,
on
=
[
'asin_bs_cate_current_id'
],
how
=
'left'
)
df
=
df
.
withColumn
(
df
=
df
.
withColumn
(
"asin_is_need"
,
"asin_is_need"
,
...
@@ -648,11 +651,11 @@ class KafkaFlowAsinDetail(Templates):
...
@@ -648,11 +651,11 @@ class KafkaFlowAsinDetail(Templates):
)
)
df
=
df
.
withColumn
(
df
=
df
.
withColumn
(
"asin_type"
,
"asin_type"
,
F
.
when
(
F
.
col
(
"asin_is_self"
)
==
1
,
F
.
lit
(
1
))
F
.
when
(
(
F
.
col
(
"asin_is_self"
)
==
1
)
|
(
F
.
col
(
"asin_is_self_brand"
)
==
1
)
,
F
.
lit
(
1
))
.
when
(
F
.
col
(
"asin_is_need"
)
==
1
,
F
.
lit
(
2
))
.
when
(
F
.
col
(
"asin_is_need"
)
==
1
,
F
.
lit
(
2
))
.
when
(
F
.
col
(
"hide_flag"
)
==
1
,
F
.
lit
(
3
))
.
when
(
F
.
col
(
"hide_flag"
)
==
1
,
F
.
lit
(
3
))
.
otherwise
(
F
.
lit
(
0
))
.
cast
(
"tinyint"
)
.
otherwise
(
F
.
lit
(
0
))
.
cast
(
"tinyint"
)
)
.
drop
(
"asin_is_self"
,
"asin_is_need"
,
"hide_flag"
)
)
.
drop
(
"asin_is_self"
,
"asin_is_
self_brand"
,
"asin_is_
need"
,
"hide_flag"
)
return
df
return
df
# 12. 处理变化率相关字段(环比_mom / 同比_yoy 统一处理)
# 12. 处理变化率相关字段(环比_mom / 同比_yoy 统一处理)
...
@@ -965,11 +968,18 @@ class KafkaFlowAsinDetail(Templates):
...
@@ -965,11 +968,18 @@ class KafkaFlowAsinDetail(Templates):
print
(
"sql="
,
sql
)
print
(
"sql="
,
sql
)
self
.
df_self_asin
=
self
.
spark
.
sql
(
sqlQuery
=
sql
)
.
dropDuplicates
([
'asin'
])
.
persist
(
StorageLevel
.
DISK_ONLY
)
self
.
df_self_asin
=
self
.
spark
.
sql
(
sqlQuery
=
sql
)
.
dropDuplicates
([
'asin'
])
.
persist
(
StorageLevel
.
DISK_ONLY
)
self
.
df_self_asin
.
show
(
10
,
truncate
=
False
)
self
.
df_self_asin
.
show
(
10
,
truncate
=
False
)
print
(
"9b. 读取内部品牌信息(amazon_brand),用于asin_type计算"
)
mysql_con
=
DBUtil
.
get_connection_info
(
"mysql"
,
"us"
)
mysql_brand_sql
=
"""
SELECT DISTINCT LOWER(TRIM(brand_name)) AS brand, 1 as asin_is_self_brand FROM amazon_brand WHERE brand_type = '1'
"""
self
.
df_self_brand
=
F
.
broadcast
(
SparkUtil
.
read_jdbc_query
(
session
=
self
.
spark
,
url
=
mysql_con
[
'url'
],
pwd
=
mysql_con
[
'pwd'
],
username
=
mysql_con
[
'username'
],
query
=
mysql_brand_sql
))
self
.
df_self_brand
.
show
(
10
,
truncate
=
False
)
print
(
"10. 读取隐藏分类(流量选品模块,id_path前缀匹配),用于asin_type计算"
)
print
(
"10. 读取隐藏分类(流量选品模块,id_path前缀匹配),用于asin_type计算"
)
# category_full_name + category_disable_config 按 id_path 前缀匹配(对齐 dwt.handle_asin_is_hide)
# category_full_name + category_disable_config 按 id_path 前缀匹配(对齐 dwt.handle_asin_is_hide)
# 配置表集中存在 us(selection) 库,固定用 us 连接,实际站点靠 site 字段过滤
# category_id 列直接别名成 asin_bs_cate_current_id,跟这个阶段 df 里还没改名前的当前分类列名对齐
mysql_con
=
DBUtil
.
get_connection_info
(
"mysql"
,
"us"
)
sql
=
f
"""
sql
=
f
"""
SELECT DISTINCT category_id AS asin_bs_cate_current_id, 1 as hide_flag FROM category_full_name a
SELECT DISTINCT category_id AS asin_bs_cate_current_id, 1 as hide_flag FROM category_full_name a
WHERE EXISTS (
WHERE EXISTS (
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment