集成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")
相关产品推荐
相关产品推荐

