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

如何定期将点击流数据关联更新至Solr产品主索引?

Solution: Sync Clickstream Signals to Solr Product Core

Great question! Let’s walk through how to implement this regular sync between your clickstream data and product Solr core, covering both your proposed storage options (second Solr core or NoSQL database) and key best practices.

First: Align Your Data Structures

First, let’s make sure your clickstream storage schema plays nicely with your product core’s dynamic fields. Your product core already has the right dynamic fields (result_id* and ClickStream_score*), so we just need to ensure the clickstream source uses matching field patterns.


Scenario 1: Clickstream Data Stored in a Second Solr Core

Step 1: Configure the Clickstream Core Schema

Set up the second core’s schema.xml to store unanalyzed search terms and their Top10 results:

<field name="search_term" type="string" indexed="true" stored="true"/>
<dynamicField name="result_id_*" type="string" indexed="true" stored="true"/>
<dynamicField name="ClickStream_score_*" type="double" indexed="true" stored="true"/>

This matches the flat structure you provided in your example data.

Step 2: Regular Sync to Product Core

You have two reliable options here:

Option A: Use Solr DataImportHandler (DIH) for Cross-Core Joins

DIH lets you pull data from the clickstream core, join it with your product core, and update fields automatically. Here’s a sample data-config.xml for your product core:

<dataConfig>
  <!-- Define data sources for both cores -->
  <dataSource name="clickstream_ds" type="SolrDataSource" baseUrl="http://localhost:8983/solr/clickstream_core"/>
  <dataSource name="product_ds" type="SolrDataSource" baseUrl="http://localhost:8983/solr/product_core"/>

  <document>
    <!-- Pull all clickstream documents -->
    <entity name="cs" dataSource="clickstream_ds" query="*:*" processor="SolrEntityProcessor">
      <!-- Join with product docs that match any of the Top10 result IDs -->
      <entity name="product" dataSource="product_ds" 
              query="id:${cs.result_id_1} OR id:${cs.result_id_2} OR id:${cs.result_id_3} OR id:${cs.result_id_4} OR id:${cs.result_id_5} OR id:${cs.result_id_6} OR id:${cs.result_id_7} OR id:${cs.result_id_8} OR id:${cs.result_id_9} OR id:${cs.result_id_10}"
              processor="SolrEntityProcessor">
        <!-- Map clickstream fields to product core's dynamic fields -->
        <field column="result_id_1" name="result_id_1"/>
        <field column="ClickStream_score_1" name="ClickStream_score_1"/>
        <!-- Repeat for result_id_2 to result_id_10 and their scores -->
        
        <!-- Optional: Add the search term to a multi-valued field on the product doc -->
        <field column="search_term" name="associated_search_terms" multiValued="true"/>
      </entity>
    </entity>
  </document>
</dataConfig>

Then, configure a scheduled task in your product core’s solrconfig.xml (using a cron expression) to run this DIH job periodically.

Option B: Custom Script (Python/Java) for Full Control

If you need more flexibility (like handling incremental updates or custom logic), write a script to sync data manually. Here’s a Python example using requests:

import requests

# Solr endpoints
PRODUCT_CORE_URL = "http://localhost:8983/solr/product_core"
CLICKSTREAM_CORE_URL = "http://localhost:8983/solr/clickstream_core"

def fetch_clickstream_data():
    """Pull all clickstream docs (with pagination for large datasets)"""
    all_docs = []
    start = 0
    rows = 1000
    while True:
        response = requests.get(f"{CLICKSTREAM_CORE_URL}/select", params={
            "q": "*:*",
            "rows": rows,
            "start": start
        }).json()
        docs = response["response"]["docs"]
        if not docs:
            break
        all_docs.extend(docs)
        start += rows
    return all_docs

def update_product_docs(clickstream_docs):
    """Update product docs with clickstream signals"""
    batch_updates = []
    for cs_doc in clickstream_docs:
        search_term = cs_doc["search_term"]
        # Iterate through Top10 results
        for i in range(1, 11):
            result_id_key = f"result_id_{i}"
            score_key = f"ClickStream_score_{i}"
            if result_id_key not in cs_doc:
                break  # Stop if we hit a missing result (e.g., only Top2 for a term)
            
            result_id = cs_doc[result_id_key]
            score = cs_doc[score_key]
            
            # Fetch the matching product doc
            product_response = requests.get(f"{PRODUCT_CORE_URL}/select", params={
                "q": f"id:{result_id}",
                "rows": 1
            }).json()
            product_docs = product_response["response"]["docs"]
            if not product_docs:
                continue  # Skip if product doesn't exist
            
            product_doc = product_docs[0]
            # Update dynamic fields
            product_doc[result_id_key] = result_id
            product_doc[score_key] = score
            # Optional: Add search term to associated terms
            if "associated_search_terms" not in product_doc:
                product_doc["associated_search_terms"] = []
            if search_term not in product_doc["associated_search_terms"]:
                product_doc["associated_search_terms"].append(search_term)
            
            # Add to batch update
            batch_updates.append({
                "id": result_id,
                "doc": product_doc,
                "commitWithin": 1000  # Let Solr handle commits efficiently
            })
    
    # Submit batch updates
    if batch_updates:
        requests.post(f"{PRODUCT_CORE_URL}/update", json=batch_updates).json()

if __name__ == "__main__":
    cs_docs = fetch_clickstream_data()
    update_product_docs(cs_docs)

Schedule this script with a tool like cron (Linux) or Task Scheduler (Windows), or use a workflow tool like Airflow for more complex scheduling.


Scenario 2: Clickstream Data Stored in a NoSQL Database

Step 1: Optimize NoSQL Data Structure

For easier maintenance, consider storing Top10 results as an array instead of flat fields (though flat fields work too):

{
  "search_term": "iphone",
  "top_results": [
    {"doc_id": "3", "score": 12},
    {"doc_id": "8", "score": 234},
    {"doc_id": "9", "score": 34}
  ]
}

This makes iterating through results cleaner.

Step 2: Custom Sync Script

The workflow is similar to the Solr core scenario—just swap the data source to your NoSQL database. Here’s a MongoDB + Python example:

import pymongo
import requests

PRODUCT_CORE_URL = "http://localhost:8983/solr/product_core"
# MongoDB setup
mongo_client = pymongo.MongoClient("mongodb://localhost:27017/")
db = mongo_client["clickstream_db"]
collection = db["clickstream_data"]

def fetch_clickstream_data():
    return list(collection.find({}))

def update_product_docs(clickstream_docs):
    batch_updates = []
    for cs_doc in clickstream_docs:
        search_term = cs_doc["search_term"]
        # Iterate through array of top results
        for idx, result in enumerate(cs_doc["top_results"], 1):
            result_id = result["doc_id"]
            score = result["score"]
            result_id_key = f"result_id_{idx}"
            score_key = f"ClickStream_score_{idx}"
            
            # Fetch product doc
            product_response = requests.get(f"{PRODUCT_CORE_URL}/select", params={
                "q": f"id:{result_id}",
                "rows": 1
            }).json()
            product_docs = product_response["response"]["docs"]
            if not product_docs:
                continue
            
            product_doc = product_docs[0]
            product_doc[result_id_key] = result_id
            product_doc[score_key] = score
            # Add associated search term
            if "associated_search_terms" not in product_doc:
                product_doc["associated_search_terms"] = []
            if search_term not in product_doc["associated_search_terms"]:
                product_doc["associated_search_terms"].append(search_term)
            
            batch_updates.append({
                "id": result_id,
                "doc": product_doc,
                "commitWithin": 1000
            })
    
    if batch_updates:
        requests.post(f"{PRODUCT_CORE_URL}/update", json=batch_updates).json()

if __name__ == "__main__":
    cs_docs = fetch_clickstream_data()
    update_product_docs(cs_docs)

Key Best Practices

  • Avoid Field Conflicts: If a product doc might appear in Top10 results for multiple search terms, consider renaming dynamic fields to include the search term (e.g., result_id_iphone_1 instead of result_id_1). Your existing dynamic field pattern (result_id*) supports this.
  • Batch Updates: Always use batch updates instead of single-document updates to avoid overwhelming Solr.
  • Incremental Sync: For large datasets, add a timestamp field to clickstream docs and only sync data modified since the last run—this reduces sync time and resource usage.
  • Commit Strategy: Use commitWithin instead of manual commits to let Solr optimize write performance.

内容的提问来源于stack exchange,提问作者NAVNEET MATHPAL

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:02:52