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

集成Kafka与Flink的Drools规则引擎无法工作:找不到pom.properties

问题描述

我正在搭建一个从Kafka数据流获取输入的Drools规则引擎,Kafka到Flink的数据流运行正常,代码能接收数据,但Drools引擎无法处理数据。错误日志显示无法找到pom.properties、无法回退到pom.xml、无法构建kmodule.xml索引,最终可用KieBases为空。尝试过移动pom.properties文件、手动设置ReleaseId等方案,问题仍未解决,寻求帮助。

错误日志
2025-03-17 14:51:54 kmodule.xml found and loaded.
2025-03-17 14:51:54 2025-03-17 03:51:54,053 INFO  org.drools.compiler.kie.builder.impl.ClasspathKieProject     [] - Found kmodule: jar:file:/tmp/tm_172.18.0.5:40077-1a373e/blobStorage/job_79120341f1c31c833ed62c9081009d88/blob_p-50d87c38beef5c86faf3007e80f59dc2b8ac2fab-799229e26bc6c7811d686b3d3a63fdc4!/META-INF/kmodule.xml
2025-03-17 14:51:54 2025-03-17 03:51:54,099 WARN  org.drools.compiler.kie.builder.impl.ClasspathKieProject     [] - Unable to find pom.properties in 40077-1a373e/blobStorage/job_79120341f1c31c833ed62c9081009d88/blob_p-50d87c38beef5c86faf3007e80f59dc2b8ac2fab-799229e26bc6c7811d686b3d3a63fdc4
2025-03-17 14:51:54 2025-03-17 03:51:54,099 WARN  org.drools.compiler.kie.builder.impl.ClasspathKieProject     [] - As folder project tried to fall back to pom.xml, but could not find one
2025-03-17 14:51:54 2025-03-17 03:51:54,099 WARN  org.drools.compiler.kie.builder.impl.ClasspathKieProject     [] - Unable to load pom.properties from/tmp/tm_172.18.0.5:40077-1a373e/blobStorage/job_79120341f1c31c833ed62c9081009d88/blob_p-50d87c38beef5c86faf3007e80f59dc2b8ac2fab-799229e26bc6c7811d686b3d3a63fdc4
2025-03-17 14:51:54 2025-03-17 03:51:54,099 WARN  org.drools.compiler.kie.builder.impl.ClasspathKieProject     [] - Cannot find maven pom properties for this project. Using the container's default ReleaseId
2025-03-17 14:51:54 2025-03-17 03:51:54,100 INFO  org.drools.compiler.kie.builder.impl.InternalKieModuleProvider [] - Creating KieModule for artifact org.default:artifact:1.0.0
2025-03-17 14:51:54 2025-03-17 03:51:54,101 ERROR org.drools.compiler.kie.builder.impl.ClasspathKieProject     [] - Unable to build index of kmodule.xml url=jar:file:/tmp/tm_172.18.0.5:40077-1a373e/blobStorage/job_79120341f1c31c833ed62c9081009d88/blob_p-50d87c38beef5c86faf3007e80f59dc2b8ac2fab-799229e26bc6c7811d686b3d3a63fdc4!/META-INF/kmodule.xml
2025-03-17 14:51:54 Unable to get all ZipFile entries: 40077-1a373e/blobStorage/job_79120341f1c31c833ed62c9081009d88/blob_p-50d87c38beef5c86faf3007e80f59dc2b8ac2fab-799229e26bc6c7811d686b3d3a63fdc4
2025-03-17 14:51:54 kContainer []
2025-03-17 14:51:54 Available KieBases: []
相关配置文件及代码

pom.xml

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
  xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
  <modelVersion>4.0.0</modelVersion>

  <groupId>com.example</groupId>
  <artifactId>transactionruleengine</artifactId>
  <version>1.0-SNAPSHOT</version>
  <packaging>jar</packaging>

  <name>transactionruleengine</name>
  <url>http://maven.apache.org</url>

  <properties>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    <drools.version>8.44.0.Final</drools.version>
  </properties>

  <dependencies>
    <dependency>
      <groupId>junit</groupId>
      <artifactId>junit</artifactId>
      <version>3.8.1</version>
      <scope>test</scope>
    </dependency>


    <!-- Drools -->
    <dependency>
      <groupId>org.drools</groupId>
      <artifactId>drools-core</artifactId>
      <version>${drools.version}</version>
    </dependency>

    <dependency>
      <groupId>org.drools</groupId>
      <artifactId>drools-compiler</artifactId>
      <version>${drools.version}</version>
    </dependency>

    <dependency>
      <groupId>org.drools</groupId>
      <artifactId>drools-mvel</artifactId>
      <version>${drools.version}</version>
    </dependency>

    <dependency>
      <groupId>org.drools</groupId>
      <artifactId>drools-io</artifactId>
      <version>${drools.version}</version>
    </dependency>

    <dependency>
      <groupId>org.drools</groupId>
      <artifactId>drools-xml-support</artifactId>
      <version>${drools.version}</version>
    </dependency>

    <dependency>
      <groupId>org.kie</groupId>
      <artifactId>kie-api</artifactId>
      <version>${drools.version}</version>
    </dependency>



  </dependencies>

  <build>
    <plugins>
      <!-- Shade plugin to build a fat JAR -->
      <plugin>
        <groupId>org.apache.maven.plugins</groupId>
        <artifactId>maven-jar-plugin</artifactId>
        <version>3.2.2</version>
        <executions>
          <execution>
            <phase>package</phase>
            <configuration>
              <archive>
                <addMavenDescriptor>true</addMavenDescriptor>
              </archive>
              <filters>
                <filter>
                  <artifact>*:*</artifact>
                  <excludes>
                    <exclude>META-INF/*.SF</exclude>
                    <exclude>META-INF/*.DSA</exclude>
                    <exclude>META-INF/*.RSA</exclude>
                    <exclude>module-info.class</exclude>
                  </excludes>
                </filter>
              </filters>
              <transformers>
                <!-- Merge LICENSE and NOTICE files -->
                <transformer implementation="org.apache.maven.plugins.shade.resource.AppendingTransformer">
                  <resource>META-INF/LICENSE.md</resource>
                </transformer>
                <transformer implementation="org.apache.maven.plugins.shade.resource.AppendingTransformer">
                  <resource>META-INF/NOTICE.md</resource>
                </transformer>
                <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
                  <mainClass>com.example.TransactionProcessor</mainClass>
                </transformer>
              </transformers>
            </configuration>
          </execution>
        </executions>
      </plugin>
    </plugins>

    <resources>
      <resource>
        <directory>src/main/resources</directory>
        <filtering>true</filtering>
        <includes>
          <include>**/*</include>
        </includes>
      </resource>
    </resources>
  </build>



</project>

主代码(TransactionProcessor.java)

package com.example;

import com.fasterxml.jackson.databind.ObjectMapper;
import org.kie.api.KieBase;
import org.kie.api.KieServices;
import org.kie.api.builder.ReleaseId;
import org.kie.api.runtime.KieContainer;
import org.kie.api.runtime.KieSession;

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.api.common.serialization.SimpleStringSchema;

import java.io.InputStream;

public class TransactionProcessor {
    public static void main(String[] args) throws Exception {
        // Kafka source configuration
        KafkaSource<String> source = KafkaSource.<String>builder()
                .setBootstrapServers("host.docker.internal:9092")
                .setTopics("test-topic")
                .setGroupId("flink-group")
                .setStartingOffsets(OffsetsInitializer.earliest())
                .setValueOnlyDeserializer(new SimpleStringSchema())
                .build();

        // Flink execution environment setup
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // Processing Kafka messages
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source")
                .map(json -> {
                    if (json == null) {
                        System.err.println("Received null message from Kafka.");
                        return null;  // skip the null message
                    }

                    ObjectMapper mapper = new ObjectMapper();
                    Transaction transaction;

                    // Deserialize JSON to Transaction object
                    try {
                        transaction = mapper.readValue(json, Transaction.class);
                    } catch (Exception e) {
                        System.err.println("Error deserializing JSON: " + e.getMessage());
                        return null;  // skip invalid JSON
                    }

                    // Debug: Print parsed transaction
                    System.out.println("Parsed transaction: " + transaction);

                    InputStream is = TransactionProcessor.class.getClassLoader().getResourceAsStream("META-INF/kmodule.xml");
                    if (is == null) {
                        throw new RuntimeException("kmodule.xml not found! Ensure it's in src/main/resources/META-INF/");
                    } else {
                        System.out.println("kmodule.xml found and loaded.");
                    }


                    // Drools rule engine setup
                    KieServices ks = KieServices.Factory.get();
                    KieContainer kContainer = ks.getKieClasspathContainer();
                    System.out.println("kContainer " + kContainer.verify().getMessages().toString());


                    // Initialize KieSession
                    System.out.println("Available KieBases: " + kContainer.getKieBaseNames());
                    KieBase kBase = kContainer.getKieBase();
                    KieSession kSession = kBase.newKieSession();
                    if (kSession == null) {
                        throw new RuntimeException("KieSession initialization failed.");
                    }

                    // Insert transaction and fire all rules
                    kSession.insert(transaction);
                    kSession.fireAllRules();
                    kSession.dispose();

                    return transaction;
                })
                .print();  // Output result to console

        // Execute Flink job
        env.execute("Flink Kafka Integration Example");
    }
}

kmodule.xml

<?xml version="1.0" encoding="UTF-8"?>
<kmodule xmlns="http://www.drools.org/xsd/kmodule">
    <kbase name="rules" packages="com.example.rules">
        <ksession name="ksession-rules"/>
    </kbase>
</kmodule>

pom.properties

artifactId=transactionruleengine
groupId=com.example
version=1.0-SNAPSHOT

文件结构

  • src/main/java/com/example/
    • TransactionProcessor.java
    • Transaction.java
  • src/main/resources/META-INF/
    • kmodule.xml
    • pom.properties
  • pom.xml
解决方案

1. 修复Maven打包配置错误

当前pom.xml错误使用maven-jar-plugin配置shade插件逻辑,导致打包后Jar缺失元数据和资源。替换为maven-shade-plugin,确保正确构建Fat Jar:

<build>
    <plugins>
        <plugin>
            <groupId>org.apache.maven.plugins</groupId>
            <artifactId>maven-shade-plugin</artifactId>
            <version>3.4.1</version>
            <executions>
                <execution>
                    <phase>package</phase>
                    <goals>
                        <goal>shade</goal>
                    </goals>
                    <configuration>
                        <filters>
                            <filter>
                                <artifact>*:*</artifact>
                                <excludes>
                                    <exclude>META-INF/*.SF</exclude>
                                    <exclude>META-INF/*.DSA</exclude>
                                    <exclude>META-INF/*.RSA</exclude>
                                    <exclude>module-info.class</exclude>
                                </excludes>
                            </filter>
                        </filters>
                        <transformers>
                            <transformer implementation="org.apache.maven.plugins.shade.resource.AppendingTransformer">
                                <resource>META-INF/LICENSE.md</resource>
                            </transformer>
                            <transformer implementation="org.apache.maven.plugins.shade.resource.AppendingTransformer">
                                <resource>META-INF/NOTICE.md</resource>
                            </transformer>
                            <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
                                <mainClass>com.example.TransactionProcessor</mainClass>
                            </transformer>
                            <transformer implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"/>
                        </transformers>
                    </configuration>
                </execution>
            </executions>
        </plugin>
    </plugins>
    <resources>
        <resource>
            <directory>src/main/resources</directory>
            <filtering>false</filtering>
            <includes>
                <include>**/*</include>
            </includes>
        </resource>
    </resources>
</build>

2. 调整pom.properties的位置

Drools默认在META-INF/maven/{groupId}/{artifactId}/路径查找pom.properties,需创建对应目录并移动文件:

  • 创建目录:src/main/resources/META-INF/maven/com.example/transactionruleengine/
  • 将pom.properties移动到该目录下

3. 优化Drools容器初始化逻辑

在Flink的map算子中重复初始化KieContainer会影响性能且易引发资源问题,将其移到算子外部作为全局资源:

public class TransactionProcessor {
    // 全局KieContainer,仅初始化一次
    private static final KieContainer KIE_CONTAINER;

    static {
        KieServices ks = KieServices.Factory.get();
        KIE_CONTAINER = ks.getKieClasspathContainer();
        System.out.println("Available KieBases: " + KIE_CONTAINER.getKieBaseNames());
        if (KIE_CONTAINER.getKieBaseNames().isEmpty()) {
            throw new RuntimeException("No KieBases found. Check kmodule.xml and rule files.");
        }
    }

    public static void main(String[] args) throws Exception {
        // Kafka和Flink初始化逻辑不变...

        env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source")
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 23:27:14