Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion dtable_events/__init__.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
from dtable_events.db import init_db_session_class
from dtable_events.activities.db import get_table_activities, get_activities_detail
from dtable_events.app.event_redis import RedisClient
from dtable_events.app.event_redis import RedisClient, redis_cache
from dtable_events.statistics.db import get_user_activity_stats_by_day, get_daily_active_users, get_email_sending_logs
from dtable_events.dtable_io.excel import get_insert_update_rows
from dtable_events.dtable_io.utils import update_page_design_static_image, rename_universal_app_static_assets_dir, \
Expand Down
223 changes: 202 additions & 21 deletions dtable_events/activities/db.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@

from dtable_events.activities.models import Activities
from dtable_events.app.config import TIME_ZONE
from dtable_events.app.event_redis import redis_cache

logger = logging.getLogger(__name__)

Expand All @@ -30,6 +31,8 @@
]

DETAIL_LIMIT = 65535 # 2^16 - 1
TABLE_ACTIVITIES_CACHE_KEY = 'table_activities:{to_tz}:{dtable_uuid}:{date}'
TABLE_ACTIVITIES_CACHE_TTL = 24 * 60 * 60


class TableActivity(object):
Expand Down Expand Up @@ -240,6 +243,154 @@ def get_shifted_days_ago(offset_str, days):
return days_ago_start


def _get_timezone(offset_str):
match = re.match(r'([+-])(\d{1,2}):(\d{2})', offset_str)
if not match:
raise ValueError("Offset format must be like '+8:00' or '-9:00'")

sign = 1 if match.group(1) == '+' else -1

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Critical] 分钟级时区被忽略

Why this matters:
to_tz 来自前端的 dayjs().format().slice(-6),会包含如 +05:30 的有效偏移;这里仅使用小时构造 Etc/GMT,完全丢弃分钟。比如 UTC 18:45 在 +05:30 已是次日 00:15,但实现仍按 +05:00 归到前一天,因此查询范围、返回日期和 Redis key 都会落在错误的自然日,并会把错误统计缓存 24 小时。

Suggested fix: 解析并保留 hours 和 minutes,使用 pytz.FixedOffset(sign * (hours * 60 + minutes))(或等价的固定偏移 tzinfo),并补充跨本地午夜的 +05:30/负分钟偏移回归测试。

hours = int(match.group(2)) * sign
gmt_offset = f"Etc/GMT{'-' if hours >= 0 else '+'}{abs(hours)}"
return pytz.timezone(gmt_offset)


def _get_local_day_utc_range(day_start_local):
day_end_local = day_start_local + timedelta(days=1)
return day_start_local.astimezone(pytz.utc), day_end_local.astimezone(pytz.utc)


def _get_activity_cache_key(dtable_uuid, to_tz, date_str):
return TABLE_ACTIVITIES_CACHE_KEY.format(dtable_uuid=dtable_uuid, to_tz=to_tz, date=date_str)


def _serialize_cached_activity(op_date, insert_row, modify_row, delete_row):
return json.dumps({
'op_date': op_date.strftime('%Y-%m-%d %H:%M:%S') if op_date else None,
'insert_row': int(insert_row or 0),
'modify_row': int(modify_row or 0),
'delete_row': int(delete_row or 0),
})


def _deserialize_cached_activity(dtable_uuid, date_str, cached_value):
if not cached_value:
return None

try:
data = json.loads(cached_value)
table_activity = TableActivity()
table_activity.dtable_uuid = dtable_uuid
table_activity.date = date_str
table_activity.op_date = datetime.strptime(data['op_date'], '%Y-%m-%d %H:%M:%S') if data.get('op_date') else None
table_activity.insert_row = data.get('insert_row', 0)
table_activity.modify_row = data.get('modify_row', 0)
table_activity.delete_row = data.get('delete_row', 0)
return table_activity
except Exception as e:
logger.warning('Deserialize table activities cache failed: %s', e)
return None


def _query_table_activities_by_date(session, uuid_list, day_start_local, to_tz):
if not uuid_list:
return []

day_start_utc, day_end_utc = _get_local_day_utc_range(day_start_local)
date_str = day_start_local.strftime('%Y-%m-%d 00:00:00')

stmt = select(
Activities.dtable_uuid,
func.max(Activities.op_time).label('op_date'),
func.sum(case((Activities.op_type == 'insert_row', Activities.row_count))).label('insert_row'),
func.sum(case((Activities.op_type == 'modify_row', Activities.row_count))).label('modify_row'),
func.sum(case((Activities.op_type == 'delete_row', Activities.row_count))).label('delete_row')
).where(
Activities.op_time >= day_start_utc,
Activities.op_time < day_end_utc,
Activities.dtable_uuid.in_(uuid_list)
).group_by(Activities.dtable_uuid)

activities = session.execute(stmt).all()
activity_map = {
dtable_uuid: {
'op_date': op_date,
'insert_row': int(insert_row or 0),
'modify_row': int(modify_row or 0),
'delete_row': int(delete_row or 0),
}
for dtable_uuid, op_date, insert_row, modify_row, delete_row in activities
}

table_activities = list()
for dtable_uuid in uuid_list:
activity_info = activity_map.get(dtable_uuid, {
'op_date': day_start_local,
'insert_row': 0,
'modify_row': 0,
'delete_row': 0,
})

table_activity = TableActivity()
table_activity.dtable_uuid = dtable_uuid
table_activity.op_date = activity_info['op_date']
table_activity.date = date_str
table_activity.insert_row = activity_info['insert_row']
table_activity.modify_row = activity_info['modify_row']
table_activity.delete_row = activity_info['delete_row']
table_activities.append(table_activity)

cache_key = _get_activity_cache_key(dtable_uuid, to_tz, day_start_local.strftime('%Y-%m-%d'))
cache_value = _serialize_cached_activity(
activity_info['op_date'],
activity_info['insert_row'],
activity_info['modify_row'],
activity_info['delete_row'],
)
try:
redis_cache.set(cache_key, cache_value, timeout=TABLE_ACTIVITIES_CACHE_TTL)
except Exception as e:
logger.warning('Set table activities cache failed: %s', e)

return table_activities


def _get_cached_table_activities(uuid_list, day_start_local, to_tz):
cached_activities = []
if not uuid_list:
return cached_activities

date_str = day_start_local.strftime('%Y-%m-%d')
cache_key_list = [_get_activity_cache_key(dtable_uuid, to_tz, date_str) for dtable_uuid in uuid_list]
try:
cached_value_list = redis_cache.mget(cache_key_list)
except Exception as e:
logger.warning('Get table activities cache failed: %s', e)
cached_value_list = [None] * len(cache_key_list)

for dtable_uuid, cached_value in zip(uuid_list, cached_value_list):
table_activity = _deserialize_cached_activity(dtable_uuid, f'{date_str} 00:00:00', cached_value)
if table_activity:
cached_activities.append(table_activity)

return cached_activities


def _get_activity_day_ranges(days, to_tz):
target_timezone = _get_timezone(to_tz)
utc_now = datetime.utcnow().replace(tzinfo=pytz.utc)
local_now = utc_now.astimezone(target_timezone)
today_start_local = local_now.replace(hour=0, minute=0, second=0, microsecond=0)
start_day_local = (local_now - timedelta(days=days)).replace(hour=0, minute=0, second=0, microsecond=0)

day_ranges = []
current_day = start_day_local
while current_day <= today_start_local:
day_ranges.append(current_day)
current_day += timedelta(days=1)

return day_ranges, today_start_local


def filter_user_activate_tables(session, days, uuid_list, to_tz):
start_str = get_shifted_days_ago(to_tz, days).strftime('%Y-%m-%d %H:%M:%S')
# filter dtable_uuids
Expand All @@ -250,39 +401,69 @@ def filter_user_activate_tables(session, days, uuid_list, to_tz):
return dtable_uuids


def get_table_activities(session, uuid_list, days, start, limit, to_tz):
def get_table_activities(session, uuid_list, days, to_tz):
if not uuid_list:
return []

activities = list()
try:
uuid_list = filter_user_activate_tables(session, days, uuid_list, to_tz)
if not uuid_list:
return []
start_utc_str = get_shifted_days_ago(to_tz, days).astimezone(pytz.utc).strftime('%Y-%m-%d %H:%M:%S')
# query activities
stmt = select(
Activities.dtable_uuid, Activities.op_time.label('op_date'),
func.date_format(func.convert_tz(Activities.op_time, '+00:00', to_tz), '%Y-%m-%d 00:00:00').label('date'),
func.sum(case((Activities.op_type == 'insert_row', Activities.row_count))).label('insert_row'),
func.sum(case((Activities.op_type == 'modify_row', Activities.row_count))).label('modify_row'),
func.sum(case((Activities.op_type == 'delete_row', Activities.row_count))).label('delete_row')
).where(
Activities.op_time > start_utc_str, Activities.dtable_uuid.in_(uuid_list)
).group_by(Activities.dtable_uuid, 'date').order_by(desc(Activities.op_time)).slice(start, start + limit)
activities = session.execute(stmt).all()

day_ranges, today_start_local = _get_activity_day_ranges(days, to_tz)
activities = []
for day_start_local in day_ranges:
if day_start_local == today_start_local:
day_start_utc, day_end_utc = _get_local_day_utc_range(day_start_local)
stmt = select(
Activities.dtable_uuid,
func.max(Activities.op_time).label('op_date'),
func.sum(case((Activities.op_type == 'insert_row', Activities.row_count))).label('insert_row'),
func.sum(case((Activities.op_type == 'modify_row', Activities.row_count))).label('modify_row'),
func.sum(case((Activities.op_type == 'delete_row', Activities.row_count))).label('delete_row')
).where(
Activities.op_time >= day_start_utc,
Activities.op_time < day_end_utc,
Activities.dtable_uuid.in_(uuid_list)
).group_by(Activities.dtable_uuid)
today_activities = []
for dtable_uuid, op_date, insert_row, modify_row, delete_row in session.execute(stmt).all():
insert_row = int(insert_row or 0)
modify_row = int(modify_row or 0)
delete_row = int(delete_row or 0)
table_activity = TableActivity()
table_activity.dtable_uuid = dtable_uuid
table_activity.op_date = op_date
table_activity.date = day_start_local.strftime('%Y-%m-%d 00:00:00')
table_activity.insert_row = insert_row
table_activity.modify_row = modify_row
table_activity.delete_row = delete_row
today_activities.append(table_activity)
activities.extend(today_activities)
else:
cached_activities = _get_cached_table_activities(uuid_list, day_start_local, to_tz)
cached_uuid_set = {activity.dtable_uuid for activity in cached_activities}
activities.extend(cached_activities)

missed_uuid_list = [dtable_uuid for dtable_uuid in uuid_list if dtable_uuid not in cached_uuid_set]
if missed_uuid_list:
activities.extend(_query_table_activities_by_date(session, missed_uuid_list, day_start_local, to_tz))

activities.sort(key=lambda activity: activity.op_date or datetime.min, reverse=True)
except Exception as e:
logger.exception('Get table activities failed: %s', e)

table_activities = list()
for dtable_uuid, op_date, date, insert_row, modify_row, delete_row in activities:
for activity in activities:
if activity.insert_row == activity.modify_row == activity.delete_row == 0:
continue
table_activity = TableActivity()
table_activity.dtable_uuid = dtable_uuid
table_activity.op_date = op_date
table_activity.date = date
table_activity.insert_row = insert_row
table_activity.modify_row = modify_row
table_activity.delete_row = delete_row
table_activity.dtable_uuid = activity.dtable_uuid
table_activity.op_date = activity.op_date
table_activity.date = activity.date
table_activity.insert_row = activity.insert_row
table_activity.modify_row = activity.modify_row
table_activity.delete_row = activity.delete_row
table_activities.append(table_activity)

return table_activities
Expand Down
6 changes: 6 additions & 0 deletions dtable_events/app/event_redis.py
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,9 @@ def refresh_subscriber(self, subscriber, pubsub_channel_name, reason='unknown'):
def get(self, key):
return self._redis.get(key)

def mget(self, keys):
return self._redis.mget(keys)

def set(self, key, value, timeout=None):
if not timeout:
return self._redis.set(key, value)
Expand Down Expand Up @@ -119,6 +122,9 @@ def init_redis(self):
def get(self, key):
return self._redis_client.get(key)

def mget(self, keys):
return self._redis_client.mget(keys)

def set(self, key, value, timeout=None):
return self._redis_client.set(key, value, timeout=timeout)

Expand Down
Loading