You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Flask后端API的定时通知轮询功能最优实现方案咨询

Hey there! Let's work through the best way to handle this notification workflow in your Flask app—you’re totally right to be cautious about both spawning too many threads manually and overengineering with something heavy like Celery. Here are two practical, balanced solutions tailored to your needs:

1. APScheduler: Lightweight, Integrated Task Scheduling

APScheduler is a fantastic middle ground for Flask apps that need scheduled tasks without the overhead of a full message queue system. It’s easy to integrate, supports multiple scheduling types (one-off, interval, cron), and runs in a background thread without cluttering your app with unmanageable threads.

How to implement it:

First, install the package:

pip install apscheduler

Then, integrate it into your Flask app. Note: In production, you’ll want to store your user list and response status in a database (like PostgreSQL with SQLAlchemy) or Redis instead of in-memory variables—this example uses in-memory for simplicity:

from flask import Flask
from apscheduler.schedulers.background import BackgroundScheduler
from datetime import datetime, timedelta

app = Flask(__name__)

# Replace these with database/Redis storage in production
user_list = ["user_123", "user_456", "user_789"]
user_response_tracking = {}  # Key: user ID, Value: bool (responded or not)

def send_notification_to_next_user():
    if not user_list:
        print("No more users in the list—stopping workflow.")
        return
    
    current_user = user_list[0]
    # Your actual notification logic (email, push, etc.) goes here
    print(f"Sent notification to {current_user} at {datetime.now().strftime('%H:%M:%S')}")
    user_response_tracking[current_user] = False  # Mark as unresponded

    # Schedule a check for 1 hour from now
    scheduler.add_job(
        check_user_response,
        trigger='date',
        run_date=datetime.now() + timedelta(hours=1),
        args=[current_user]
    )

def check_user_response(user_id):
    if user_id not in user_response_tracking or not user_response_tracking[user_id]:
        # User didn't respond—remove them from the list
        if user_id in user_list:
            user_list.remove(user_id)
            print(f"Removed {user_id} (no response after 1 hour)")
        # Move to the next user
        send_notification_to_next_user()
    else:
        print(f"{user_id} responded—no further action needed.")

# Initialize the scheduler and start the workflow
scheduler = BackgroundScheduler()
scheduler.add_job(send_notification_to_next_user, trigger='date', run_date=datetime.now())
scheduler.start()

# Endpoint for users to mark themselves as responded
@app.route('/respond/<user_id>')
def record_response(user_id):
    if user_id in user_response_tracking:
        user_response_tracking[user_id] = True
        return f"Response recorded for {user_id}!"
    return "User not found", 404

if __name__ == '__main__':
    app.run(debug=True)

Key notes for production:

  • Use a persistent storage backend (database/Redis) for user_list and user_response_tracking—in-memory data will be lost if your Flask app restarts.
  • For multi-worker setups, use a shared scheduler backend like Redis or SQLAlchemy so all workers can access task state.
2. Redis + RQ (Redis Queue): Lightweight Queue with Scheduling

If you need more reliability (e.g., handling task retries, distributed workers) but still want to avoid Celery’s complexity, RQ (Redis Queue) with RQ-Scheduler is a great pick. It uses Redis as a message broker, keeps tasks organized, and supports scheduled one-off jobs.

How to implement it:

First, install the required packages:

pip install rq rq-scheduler redis

Create a tasks.py file for your background tasks:

import redis
from rq import Queue
from rq_scheduler import Scheduler
from datetime import timedelta

# Connect to Redis (configure host/port/password as needed)
redis_conn = redis.Redis(host='localhost', port=6379, db=0)
task_queue = Queue(connection=redis_conn)
scheduler = Scheduler(connection=redis_conn)

# Redis keys for persistent storage
USER_LIST_KEY = "notification_user_list"
USER_RESPONSE_KEY = "user_response:{user_id}"

def send_notification_and_schedule_check():
    # Get the first user from the Redis list
    current_user_bytes = redis_conn.lindex(USER_LIST_KEY, 0)
    if not current_user_bytes:
        print("No users left to notify.")
        return
    
    current_user = current_user_bytes.decode('utf-8')
    # Your notification logic here
    print(f"Sent notification to {current_user}")
    # Mark user as unresponded in Redis
    redis_conn.set(USER_RESPONSE_KEY.format(user_id=current_user), "false")

    # Schedule the response check for 1 hour later
    scheduler.enqueue_in(
        timedelta(hours=1),
        check_user_response,
        current_user
    )

def check_user_response(user_id):
    response_status = redis_conn.get(USER_RESPONSE_KEY.format(user_id=user_id))
    if not response_status or response_status.decode('utf-8') == "false":
        # Remove unresponsive user from the list
        redis_conn.lrem(USER_LIST_KEY, 0, user_id)
        print(f"Removed {user_id} (no response)")
        # Notify the next user
        send_notification_and_schedule_check()
    else:
        print(f"{user_id} responded successfully.")

Then, update your Flask app to trigger the workflow and handle responses:

from flask import Flask
from tasks import send_notification_and_schedule_check, redis_conn, USER_RESPONSE_KEY

app = Flask(__name__)

# Endpoint to start the notification workflow
@app.route('/start-notifications')
def start_workflow():
    # Populate the user list in Redis (replace with your actual data source)
    redis_conn.rpush("notification_user_list", "user_123", "user_456", "user_789")
    send_notification_and_schedule_check()
    return "Notification workflow started!"

# Endpoint for users to respond
@app.route('/respond/<user_id>')
def record_response(user_id):
    redis_conn.set(USER_RESPONSE_KEY.format(user_id=user_id), "true")
    return f"Thanks for responding, {user_id}!"

if __name__ == '__main__':
    app.run(debug=True)

Key notes for production:

  • Run the RQ worker and scheduler separately: rq worker and rqscheduler commands.
  • Redis ensures task and user state persists across app restarts and worker instances.
Final Recommendations
  • Avoid manual threading: Spawning a thread per user is messy—you’ll end up with unmanageable thread counts, no built-in task persistence, and potential memory leaks.
  • Skip Celery unless you need it: Celery is powerful but adds significant configuration overhead (broker, result backend, workers) that’s overkill for your straightforward workflow.
  • Choose APScheduler if: You have a single Flask instance, want minimal setup, and don’t need distributed workers.
  • Choose RQ + Redis if: You need task persistence, retries, or plan to scale to multiple workers in the future—it’s lightweight but more robust than APScheduler for distributed setups.

内容的提问来源于stack exchange,提问作者John Fisher

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 10:40:44