当前位置:网站首页 >  教程

流量投流矩阵:从零搭建自动化流量投放系统

时间:2026年05月29日 04:17:41 来源:易频IT社区

一、系统架构与核心组件

流量投流矩阵是一个自动化管理多平台广告投放的系统,通过统一接口控制多个广告平台的投放动作。核心架构分为三层:数据采集层、策略决策层和执行控制层。

1.1 技术栈选型

  • 后端框架:Python 3.9 + FastAPI
  • 数据库:PostgreSQL 14(存储投放数据)+ Redis 7(缓存实时指标)
  • 任务队列:Celery 5.3 + RabbitMQ 3.11
  • 前端监控面板:Vue 3 + ECharts 5

1.2 目录结构

创建项目目录:

mkdir traffic-matrix && cd traffic-matrix
mkdir -p src/{core,platforms,strategies,models,utils}
mkdir config logs data

二、环境配置与依赖安装

2.1 Python环境配置

创建虚拟环境并安装核心依赖:

python -m venv venv
source venv/bin/activate   Linux/Mac
venv\Scripts\activate   Windows
pip install fastapi==0.104.1 uvicorn==0.24.0
pip install sqlalchemy==2.0.23 psycopg2-binary==2.9.9 redis==5.0.1
pip install celery==5.3.4 requests==2.31.0 pandas==2.1.3

2.2 数据库初始化

创建PostgreSQL数据库:

sudo -u postgres psql
CREATE DATABASE traffic_matrix;
CREATE USER matrix_user WITH PASSWORD 'YourSecurePassword123!';
GRANT ALL PRIVILEGES ON DATABASE traffic_matrix TO matrix_user;

创建数据库连接配置文件 config/database.py

DATABASE_CONFIG = {
'host': 'localhost',
'port': 5432,
'database': 'traffic_matrix',
'user': 'matrix_user',
'password': 'YourSecurePassword123!',
'pool_size': 20,
'max_overflow': 30
}
REDIS_CONFIG = {
'host': 'localhost',
'port': 6379,
'db': 0,
'password': '',
'decode_responses': True
}

三、数据模型设计

3.1 核心表结构

创建 src/models/base.py

from sqlalchemy import create_engine, Column, Integer, String, Float, DateTime, JSON, Boolean
from sqlalchemy.ext.declarative import declarative_base
from sqlalchemy.orm import sessionmaker
from datetime import datetime
import pytz
engine = create_engine(
f"postgresql://{DATABASE_CONFIG['user']}:{DATABASE_CONFIG['password']}"
f"@{DATABASE_CONFIG['host']}:{DATABASE_CONFIG['port']}"
f"/{DATABASE_CONFIG['database']}",
pool_size=DATABASE_CONFIG['pool_size'],
max_overflow=DATABASE_CONFIG['max_overflow']
)
SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine)
Base = declarative_base()
class Campaign(Base):
__tablename__ = 'campaigns'
id = Column(Integer, primary_key=True)
platform = Column(String(50), nullable=False)   平台名称:douyin, kuaishou, etc.
campaign_id = Column(String(100), unique=True, nullable=False)
name = Column(String(200))
daily_budget = Column(Float, default=0.0)
status = Column(String(20), default='ACTIVE')
targeting = Column(JSON)   定向条件
creatives = Column(JSON)   创意素材
metrics = Column(JSON)   实时指标
created_at = Column(DateTime, default=datetime.now(pytz.UTC))
updated_at = Column(DateTime, default=datetime.now(pytz.UTC), onupdate=datetime.now(pytz.UTC))
class AdGroup(Base):
__tablename__ = 'ad_groups'
id = Column(Integer, primary_key=True)
campaign_id = Column(Integer, nullable=False)
group_id = Column(String(100), unique=True, nullable=False)
bid_amount = Column(Float, default=0.0)
optimization_goal = Column(String(50))   优化目标
placement = Column(JSON)   广告位
schedule = Column(JSON)   投放时段
status = Column(String(20), default='ACTIVE')
class PerformanceLog(Base):
__tablename__ = 'performance_logs'
id = Column(Integer, primary_key=True)
campaign_id = Column(Integer, nullable=False)
timestamp = Column(DateTime, nullable=False)
impressions = Column(Integer, default=0)
clicks = Column(Integer, default=0)
conversions = Column(Integer, default=0)
spend = Column(Float, default=0.0)
ctr = Column(Float, default=0.0)
cvr = Column(Float, default=0.0)
cpc = Column(Float, default=0.0)
roas = Column(Float, default=0.0)

初始化数据库表:

流量投流矩阵:从零搭建自动化流量投放系统

 src/init_db.py
from models.base import Base, engine
Base.metadata.create_all(bind=engine)
print("数据库表创建完成")

四、平台接口封装

4.1 抽象基类设计

创建 src/platforms/base.py

from abc import ABC, abstractmethod
import requests
import json
from datetime import datetime
import hashlib
import hmac
class PlatformAdapter(ABC):
def __init__(self, access_token, advertiser_id):
self.access_token = access_token
self.advertiser_id = advertiser_id
self.base_url = self.get_base_url()
@abstractmethod
def get_base_url(self):
"""返回平台API基础地址"""
pass
@abstractmethod
def create_campaign(self, campaign_data):
"""创建广告计划"""
pass
@abstractmethod
def update_campaign_budget(self, campaign_id, budget):
"""更新广告计划预算"""
pass
@abstractmethod
def get_campaign_metrics(self, campaign_id, start_date, end_date):
"""获取广告计划指标"""
pass
def make_request(self, method, endpoint, params=None, data=None):
"""统一请求方法"""
url = f"{self.base_url}{endpoint}"
headers = {
'Content-Type': 'application/json',
'Access-Token': self.access_token
}
response = requests.request(
method=method,
url=url,
headers=headers,
params=params,
json=data
)
if response.status_code == 200:
return response.json()
else:
raise Exception(f"API请求失败: {response.status_code}, {response.text}")

4.2 抖音平台实现

创建 src/platforms/douyin.py

from .base import PlatformAdapter
import time
class DouyinAdapter(PlatformAdapter):
def get_base_url(self):
return "https://ad.oceanengine.com/open_api/v2.0"
def create_campaign(self, campaign_data):
endpoint = "/campaign/create/"
payload = {
"advertiser_id": self.advertiser_id,
"campaign_name": campaign_data['name'],
"budget_mode": campaign_data.get('budget_mode', 'BUDGET_MODE_DAY'),
"budget": campaign_data.get('daily_budget'),
"campaign_type": campaign_data.get('campaign_type', 'VIDEO_AND_IMAGE'),
"operation": "enable"
}
return self.make_request('POST', endpoint, data=payload)
def update_campaign_budget(self, campaign_id, budget):
endpoint = "/campaign/update/"
payload = {
"advertiser_id": self.advertiser_id,
"campaign_id": campaign_id,
"budget": budget
}
return self.make_request('POST', endpoint, data=payload)
def get_campaign_metrics(self, campaign_id, start_date, end_date):
endpoint = "/report/campaign/get/"
params = {
"advertiser_id": self.advertiser_id,
"start_date": start_date,
"end_date": end_date,
"campaign_ids": [campaign_id],
"group_by": ["STAT_GROUP_BY_CAMPAIGN_ID"],
"page": 1,
"page_size": 1000
}
return self.make_request('GET', endpoint, params=params)

4.3 快手平台实现

创建 src/platforms/kuaishou.py

from .base import PlatformAdapter
import hashlib
import urllib.parse
class KuaishouAdapter(PlatformAdapter):
def get_base_url(self):
return "https://api.e.kuaishou.com/rest/openapi"
def _generate_signature(self, params):
"""生成快手API签名"""
sorted_params = sorted(params.items())
param_str = '&'.join([f'{k}={v}' for k, v in sorted_params])
sign_str = f"{param_str}&access_token={self.access_token}"
return hashlib.md5(sign_str.encode()).hexdigest()
def create_campaign(self, campaign_data):
endpoint = "/campaign/create"
params = {
"advertiser_id": self.advertiser_id,
"campaign_name": campaign_data['name'],
"day_budget": int(campaign_data.get('daily_budget', 0)  100),   转换为分
"operation": "enable",
"timestamp": int(time.time()  1000)
}
params['sign'] = self._generate_signature(params)
return self.make_request('POST', endpoint, data=params)

五、策略引擎实现

5.1 预算调整策略

创建 src/strategies/budget_strategy.py

class BudgetOptimizationStrategy:
def __init__(self, min_budget=100, max_budget=50000, step_size=0.2):
self.min_budget = min_budget
self.max_budget = max_budget
self.step_size = step_size   每次调整幅度
def calculate_new_budget(self, current_budget, roas, target_roas=2.0):
"""根据ROAS调整预算"""
if roas == 0:
return current_budget  0.8   无转化时降低预算
roas_ratio = roas / target_roas
if roas_ratio > 1.2:   ROAS超过目标20%
increase = current_budget  self.step_size
new_budget = current_budget + increase
elif roas_ratio < 0.8:   ROAS低于目标20%
decrease = current_budget  self.step_size
new_budget = current_budget - decrease
else:
new_budget = current_budget
限制预算范围
new_budget = max(self.min_budget, min(new_budget, self.max_budget))
return round(new_budget, 2)

5.2 出价调整策略

创建 src/strategies/bid_strategy.py

class BidOptimizationStrategy:
def __init__(self, min_bid=0.1, max_bid=100):
self.min_bid = min_bid
self.max_bid = max_bid
def adjust_bid_by_ctr(self, current_bid, ctr, avg_ctr):
"""根据CTR调整出价"""
if avg_ctr == 0:
return current_budget
ctr_ratio = ctr / avg_ctr
if ctr_ratio > 1.3:   CTR高于平均30%
new_bid = current_bid  1.15
elif ctr_ratio > 1.1:   CTR高于平均10%
new_bid = current_bid  1.05
elif ctr_ratio < 0.7:   CTR低于平均30%
new_bid = current_bid  0.85
else:
new_bid = current_bid
return round(max(self.min_bid, min(new_bid, self.max_bid)), 2)

六、任务调度系统

6.1 Celery配置

创建 src/celery_app.py

from celery import Celery
from config.database import REDIS_CONFIG
celery_app = Celery(
'traffic_matrix',
broker=f"redis://{REDIS_CONFIG['host']}:{REDIS_CONFIG['port']}/{REDIS_CONFIG['db']}",
backend=f"redis://{REDIS_CONFIG['host']}:{REDIS_CONFIG['port']}/{REDIS_CONFIG['db']}"
)
celery_app.conf.update(
task_serializer='json',
accept_content=['json'],
result_serializer='json',
timezone='Asia/Shanghai',
enable_utc=True,
beat_schedule={
'sync-campaign-metrics-every-30-min': {
'task': 'src.tasks.sync_metrics',
'schedule': 1800.0,   30分钟
},
'optimize-budget-every-2-hours': {
'task': 'src.tasks.optimize_budget',
'schedule': 7200.0,   2小时
},
'adjust-bids-every-hour': {
'task': 'src.tasks.adjust_bids',
'schedule': 3600.0,   1小时
}
}
)

6.2 核心定时任务

创建 src/tasks.py

from celery_app import celery_app
from models.base import SessionLocal
from models.campaign import Campaign
from datetime import datetime, timedelta
import pytz
@celery_app.task
def sync_metrics():
"""同步各平台广告数据"""
db = SessionLocal()
try:
campaigns = db.query(Campaign).filter(Campaign.status == 'ACTIVE').all()
for campaign in campaigns:
根据平台选择适配器
if campaign.platform == 'douyin':
from platforms.douyin import DouyinAdapter
adapter = DouyinAdapter(
access_token=campaign.access_token,
advertiser_id=campaign.advertiser_id
)
elif campaign.platform == 'kuaishou':
from platforms.kuaishou import KuaishouAdapter
adapter = KuaishouAdapter(
access_token=campaign.access_token,
advertiser_id=campaign.advertiser_id
)
else:
continue
获取当天数据
end_date = datetime.now(pytz.UTC)
start_date = end_date - timedelta(days=1)
metrics = adapter.get_campaign_metrics(
campaign.campaign_id,
start_date.strftime('%Y-%m-%d'),
end_date.strftime('%Y-%m-%d')
)
更新数据库
campaign.metrics = metrics
campaign.updated_at = datetime.now(pytz.UTC)
db.commit()
finally:
db.close()
@celery_app.task
def optimize_budget():
"""优化广告预算"""
db = SessionLocal()
try:
from strategies.budget_strategy import BudgetOptimizationStrategy
budget_strategy = BudgetOptimizationStrategy()
campaigns = db.query(Campaign).filter(Campaign.status == 'ACTIVE').all()
for campaign in campaigns:
if not campaign.metrics:
continue
roas = campaign.metrics.get('roas', 0)
current_budget = campaign.daily_budget
new_budget = budget_strategy.calculate_new_budget(
current_budget, roas
)
if abs(new_budget - current_budget) > current_budget  0.1:
更新平台预算
if campaign.platform == 'douyin':
from platforms.douyin import DouyinAdapter
adapter = DouyinAdapter(
access_token=campaign.access_token,
advertiser_id=campaign.advertiser_id
)
adapter.update_campaign_budget(campaign.campaign_id, new_budget)
更新数据库
campaign.daily_budget = new_budget
campaign.updated_at = datetime.now(pytz.UTC)
db.commit()
finally:
db.close()

七、API接口实现

7.1 FastAPI主应用

创建 src/main.py

from fastapi import FastAPI, Depends, HTTPException
from sqlalchemy.orm import Session
from models.base import

相关推荐

最新

热门

推荐

精选

标签

易频IT社区是综合性互联网IT技术门户网站,专注分享网络技术、服务器运维、网络安全、编程开发、系统架构、云计算、大数据等行业干货,实时更新IT行业资讯、零基础教程、实战案例,为IT从业者、技术爱好者提供专业的学习交流平台。

Copyright © 2021-2026 易频IT社区. All Rights Reserved. 备案号:闽ICP备2023013482号 网站地图