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

NiFi中Python 2.X环境下json.loads处理特殊字符(如Ñ、é等)的问题

解决Jython InvokeScriptedProcessor中特殊字符乱码问题

我之前也遇到过类似的编码坑,咱们一步步拆解问题并解决:

问题根源分析

  1. 错误的编码时机:你在拼接SQL时对字符串调用.encode("utf-8"),在Jython中这会把Unicode字符串转成字节数组。直接把字节数组拼进SQL字符串时,Jython会用系统默认编码(而非UTF-8)转成字符串,导致Ñ变成ÃUECO这类乱码。
  2. JSON序列化的默认转义:json.dumps()默认开启ensure_ascii=True,会把非ASCII字符转义成\uXXXX格式(比如Ñ变成\u00d1),不符合你的需求。

具体解决方案

1. 移除拼接SQL时的.encode("utf-8")

直接使用Unicode字符串拼接SQL,只在最后写入输出流时统一编码为UTF-8即可。

2. 修改json.dumps()参数,保留原字符

添加ensure_ascii=False参数,让JSON序列化时保留特殊字符,不进行转义。

3. 处理SQL单引号转义(必加,避免语法错误)

为了防止SQL语法报错和注入风险,需要把字符串中的'替换为''。

修改后的代码示例

调整__generate_sql_transaction函数

def __generate_sql_transaction(input_data):
    """ Generate SQL statement """
    sql = """ BEGIN;"""
    # 转义单引号,避免SQL语法错误
    _id = input_data.get("id").replace("'", "''") if input_data.get("id") else ""
    _timestamp = input_data.get("timestamp").replace("'", "''") if input_data.get("timestamp") else ""
    _flowfile_metrics = input_data.get("metrics")
    _flowfile_metadata = input_data.get("metadata")
    self.valid = __validate_metrics_type(_flowfile_metrics)
    
    if self.valid is True:
        self.log.error("generate insert")
        sql += """ INSERT INTO {0}.{1} (id, timestamp, metrics""".format(schema, table)
        if _flowfile_metadata:
            sql += ", metadata"
        
        # 关闭JSON转义,直接使用Unicode字符串
        metrics_json = json.dumps(_flowfile_metrics, ensure_ascii=False).replace("'", "''")
        sql += """) VALUES ('{0}', '{1}', '{2}'""".format(_id, _timestamp, metrics_json)
        
        self.log.error("generate metadata")
        if _flowfile_metadata:
            metadata_json = json.dumps(_flowfile_metadata, ensure_ascii=False).replace("'", "''")
            sql += ", '{}'".format(metadata_json)
        
        sql += """) ON CONFLICT ({})""".format(on_conflict)
        if not bool(int(self.update)):
            sql += """ DO NOTHING;"""
        else:
            sql += """ DO UPDATE SET"""
            metrics_update = json.dumps(_flowfile_metrics, ensure_ascii=False).replace("'", "''")
            if bool(int(self.preference)):
                sql += """ metrics = '{2}' || {0}.{1}.metrics;""".format(schema, table, metrics_update)
            else:
                sql += """ metrics = {0}.{1}.metrics || '{2}';""".format(schema, table, metrics_update)
    else:
        return ""
    
    sql += """ COMMIT;"""
    return sql

调整输出写入逻辑(原代码基础上优化)

output = __generate_sql_transaction(data)
self.log.error("post generate_sql_transaction")
self.log.error(output)  # 直接打印Unicode字符串,无需提前编码
# If no sql_transaction is generated because requisites weren't met,
# set the processor output with the original flowfile input.
if output == "":
    output = text
# 最后统一编码为UTF-8写入流
outputStream.write(output.encode("utf-8"))

效果验证

修改后生成的SQL会和你期望的一致:

BEGIN; INSERT INTO Table (id, timestamp, metrics, metadata) VALUES ('ÑUECO', '2020-01-01T00:00:00+01:00', '{"value":3.1415}', '{"location":"ÑUECO"}') ON CONFLICT (id, timestamp) DO UPDATE SET metrics='{"value":3.1415}' || Table.metrics; COMMIT;

额外提醒

手动拼接SQL存在SQL注入风险,如果你的场景允许,建议结合NiFi的PutSQL处理器使用参数化查询(通过属性传递参数),从根源避免风险。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 15:32:44