Flink CDC写入Elasticsearch 7.17.9时遇TimeValue类缺失错误
解决方案:Flink CDC写入Elasticsearch 7.17.9出现
NoClassDefFoundError: org/elasticsearch/common/unit/TimeValue 错误原因分析
该错误表明JVM运行时无法找到TimeValue类,这个类属于Elasticsearch核心客户端库,根源通常是依赖版本不匹配、依赖缺失或打包时依赖未正确包含。
步骤1:修正Maven依赖配置
确保Flink Elasticsearch Connector与ES服务端版本(7.17.9)兼容,显式引入匹配版本的ES客户端依赖,并排除冲突依赖:
<dependencies> <!-- Flink MySQL CDC 依赖 --> <dependency> <groupId>com.ververica</groupId> <artifactId>flink-connector-mysql-cdc</artifactId> <version>2.4.1</version> </dependency> <!-- Flink Elasticsearch 7 Sink 依赖 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-elasticsearch7_2.12</artifactId> <version>${flink.version}</version> <!-- 需与你的Flink版本一致,如1.17.1 --> <!-- 排除自带的ES客户端,避免版本冲突 --> <exclusions> <exclusion> <groupId>org.elasticsearch.client</groupId> <artifactId>elasticsearch-rest-high-level-client</artifactId> </exclusion> <exclusion> <groupId>org.elasticsearch</groupId> <artifactId>elasticsearch</artifactId> </exclusion> </exclusions> </dependency> <!-- 引入与ES服务端同版本的核心依赖 --> <dependency> <groupId>org.elasticsearch</groupId> <artifactId>elasticsearch</artifactId> <version>7.17.9</version> </dependency> <dependency> <groupId>org.elasticsearch.client</groupId> <artifactId>elasticsearch-rest-high-level-client</artifactId> <version>7.17.9</version> <exclusions> <exclusion> <groupId>org.elasticsearch</groupId> <artifactId>elasticsearch</artifactId> </exclusion> </exclusions> </dependency> <!-- Flink 核心依赖(scope设为provided避免打包重复) --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_2.12</artifactId> <version>${flink.version}</version> <scope>provided</scope> </dependency> </dependencies>
步骤2:确保打包时包含所有依赖
使用maven-shade-plugin构建Fat Jar,避免依赖丢失或冲突:
<build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.2.4</version> <executions> <execution> <phase>package</phase> <goals> <goal>shade</goal> </goals> <configuration> <artifactSet> <excludes> <exclude>org.apache.flink:flink-shaded-*</exclude> <exclude>org.apache.flink:flink-core</exclude> <exclude>org.apache.flink:flink-streaming-java</exclude> </excludes> </artifactSet> <transformers> <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"> <mainClass>你的主类全路径</mainClass> </transformer> <transformer implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"/> </transformers> <relocations> <relocation> <pattern>org.apache.flink</pattern> <shadedPattern>org.apache.flink.shaded</shadedPattern> </relocation> </relocations> </configuration> </execution> </executions> </plugin> </plugins> </build>
步骤3:验证ES Sink代码配置
确保使用适配ES 7.x的API,避免过时方法:
import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.connectors.elasticsearch.ElasticsearchSinkFunction; import org.apache.flink.streaming.connectors.elasticsearch.RequestIndexer; import org.apache.flink.streaming.connectors.elasticsearch7.ElasticsearchSink; import org.elasticsearch.action.index.IndexRequest; import org.elasticsearch.client.Requests; import org.apache.http.HttpHost; import java.util.ArrayList; import java.util.List; import java.util.Map; public class EsSinkHandler { public static void addElasticsearchSink(DataStream<Map<String, Object>> dataStream) { List<HttpHost> httpHosts = new ArrayList<>(); httpHosts.add(new HttpHost("localhost", 9200, "http")); ElasticsearchSink.Builder<Map<String, Object>> sinkBuilder = new ElasticsearchSink.Builder<>( httpHosts, new ElasticsearchSinkFunction<Map<String, Object>>() { @Override public void process(Map<String, Object> element, RuntimeContext ctx, RequestIndexer indexer) { IndexRequest request = Requests.indexRequest() .index("your_target_index") // 替换为你的ES索引名 .source(element); indexer.add(request); } } ); // 配置批量写入参数 sinkBuilder.setBulkFlushMaxActions(1000); sinkBuilder.setBulkFlushInterval(3000); // 3秒触发一次批量写入 dataStream.addSink(sinkBuilder.build()); } }
内容的提问来源于stack exchange,提问作者user17737891
相关产品推荐
相关产品推荐

