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

基于Parquet外部存储的Ignite内存数据网格读写穿透缓存实现求助

解决方案:基于Parquet的Ignite读写穿透缓存实现与参数配置

一、Python实现读写穿透逻辑

1. 依赖准备

  • 安装必要的Python库:
pip install pyignite pyspark

2. 核心实现代码

初始化连接

from pyignite import Client, CacheConfiguration, QueryEntity
from pyspark.sql import SparkSession

# 初始化Ignite客户端
ignite_client = Client()
ignite_client.connect('127.0.0.1', 10800)

# 初始化SparkSession
spark = SparkSession.builder \
    .appName("IgniteParquetReadWriteThrough") \
    .getOrCreate()

# 定义常量
CACHE_NAME = "parquet_backed_cache"
PARQUET_PATH = "/path/to/your/parquet/files"

# 配置缓存Schema(需与Parquet文件结构匹配)
cache_config = CacheConfiguration(
    name=CACHE_NAME,
    query_entity=[
        QueryEntity(
            type_name='DataEntity',
            key_type='java.lang.Integer',
            value_type='java.util.HashMap',
            fields=[
                QueryField('id', 'java.lang.Integer'),
                QueryField('name', 'java.lang.String'),
                QueryField('age', 'java.lang.Integer')
                # 按实际Parquet字段扩展
            ]
        )
    ]
)

# 创建或获取缓存
ignite_client.get_or_create_cache(cache_config)

读穿透查询函数

def read_through_query(query_sql):
    # 解析查询条件,兼容有无WHERE子句的情况
    if 'WHERE' in query_sql:
        base_sql, condition = query_sql.split('WHERE', 1)
        cache_query = f"{base_sql.strip()} WHERE {condition.strip()}"
        parquet_filter = condition.strip()
    else:
        cache_query = query_sql.strip()
        parquet_filter = "1=1"
    
    # 1. 优先查询Ignite缓存
    try:
        cache = ignite_client.get_cache(CACHE_NAME)
        cache_result = list(ignite_client.sql(cache_query))
        
        if cache_result:
            # 缓存命中,转换为DataFrame返回
            return spark.createDataFrame(cache_result, schema=cache.schema)
    except Exception:
        # 缓存无匹配数据或查询失败,继续查询Parquet
        pass
    
    # 2. 查询底层Parquet文件
    parquet_df = spark.read.parquet(PARQUET_PATH).filter(parquet_filter)
    
    # 3. 将查询结果写入缓存,预热缓存
    if not parquet_df.isEmpty():
        cache.put_all({row.id: row.asDict() for row in parquet_df.collect()})
    
    return parquet_df

写穿透实现

def write_through_data(df):
    # 1. 写入Ignite缓存
    cache = ignite_client.get_cache(CACHE_NAME)
    cache.put_all({row.id: row.asDict() for row in df.collect()})
    
    # 2. 同步写入Parquet文件(支持追加/覆盖模式,按需选择)
    df.write.mode("append").parquet(PARQUET_PATH)

3. 使用示例

# 读穿透查询示例
result_df = read_through_query("SELECT * FROM parquet_backed_cache WHERE age > 30")
result_df.show()

# 写穿透示例
new_data = [(101, "Alice", 32), (102, "Bob", 28)]
new_df = spark.createDataFrame(new_data, ["id", "name", "age"])
write_through_data(new_df)

二、启用Ignite的read-through/write-through参数

Ignite的read-through和write-through需结合自定义CacheStore实现才能生效,具体操作如下:

1. 通过SQL修改缓存参数

对已存在的缓存,可直接用Ignite SQL命令启用参数:

ALTER CACHE parquet_backed_cache 
SET READ_THROUGH = TRUE, 
WRITE_THROUGH = TRUE;

若创建新缓存,可直接在创建语句中指定:

CREATE CACHE parquet_backed_cache 
WITH READ_THROUGH = TRUE, 
WRITE_THROUGH = TRUE, 
QUERY_ENTITIES = [
    {
        "typeName": "DataEntity",
        "keyType": "java.lang.Integer",
        "valueType": "java.util.HashMap",
        "fields": [
            {"name": "id", "type": "java.lang.Integer"},
            {"name": "name", "type": "java.lang.String"},
            {"name": "age", "type": "java.lang.Integer"}
        ]
    }
];

2. Python实现Parquet CacheStore

基于pyignite的CacheStore抽象类,实现对接Parquet文件的存储逻辑:

from pyignite.cache_store import CacheStore

class ParquetCacheStore(CacheStore):
    def __init__(self, parquet_path):
        self.parquet_path = parquet_path
        self.spark = SparkSession.builder.appName("ParquetCacheStore").getOrCreate()
    
    def load(self, key):
        # 从Parquet加载单个key对应数据
        df = self.spark.read.parquet(self.parquet_path).filter(f"id = {key}")
        return df.collect()[0].asDict() if not df.isEmpty() else None
    
    def load_all(self, keys):
        # 批量加载数据
        keys_str = ",".join(map(str, keys))
        df = self.spark.read.parquet(self.parquet_path).filter(f"id IN ({keys_str})")
        return {row.id: row.asDict() for row in df.collect()}
    
    def write(self, key, value):
        # 写入单个数据到Parquet
        df = self.spark.createDataFrame([value])
        df.write.mode("append").parquet(self.parquet_path)
    
    def write_all(self, entries):
        # 批量写入数据到Parquet
        df = self.spark.createDataFrame(entries.values())
        df.write.mode("append").parquet(self.parquet_path)
    
    def delete(self, key):
        # 删除Parquet中对应key的数据(Parquet不支持随机删除,这里用重写方式实现)
        df = self.spark.read.parquet(self.parquet_path).filter(f"id != {key}")
        df.write.mode("overwrite").parquet(self.parquet_path)
    
    def delete_all(self, keys):
        # 批量删除数据
        keys_str = ",".join(map(str, keys))
        df = self.spark.read.parquet(self.parquet_path).filter(f"id NOT IN ({keys_str})")
        df.write.mode("overwrite").parquet(self.parquet_path)

3. 绑定CacheStore到缓存

在初始化缓存时,将自定义Store绑定到缓存配置:

cache_config = CacheConfiguration(
    name=CACHE_NAME,
    read_through=True,
    write_through=True,
    cache_store_factory=lambda: ParquetCacheStore(PARQUET_PATH),
    query_entity=[
        QueryEntity(
            type_name='DataEntity',
            key_type='java.lang.Integer',
            value_type='java.util.HashMap',
            fields=[
                QueryField('id', 'java.lang.Integer'),
                QueryField('name', 'java.lang.String'),
                QueryField('age', 'java.lang.Integer')
            ]
        )
    ]
)

ignite_client.create_cache(cache_config)

内容的提问来源于stack exchange,提问作者Arunima Barik

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 06:30:44