
1. 数据采集基础概念与工具选型数据采集作为数据处理流程的第一步其质量直接影响后续分析和应用的可靠性。在实际项目中我们通常需要根据数据源类型、数据量和实时性要求选择合适的采集工具。1.1 常见数据采集场景分类根据我的项目经验数据采集主要分为以下几种典型场景日志文件采集服务器日志、应用日志等文本数据的持续收集数据库变更捕获MySQL binlog、MongoDB oplog等数据库变更流的获取API接口调用通过HTTP/RESTful接口获取结构化数据网页爬取从网页中提取非结构化或半结构化数据物联网设备数据传感器、智能设备的时序数据采集1.2 主流数据采集工具对比针对不同场景我们有以下工具可选工具名称适用场景优点缺点典型应用Flume日志收集高可靠、支持多种source/sink配置复杂服务器日志采集Logstash日志处理插件丰富、过滤能力强资源消耗大ELK日志系统Kafka Connect数据管道高吞吐、分布式需要Kafka环境数据中台建设Python爬虫网页抓取灵活可控需要编码竞品数据采集Prometheus指标监控时序数据处理强只适合指标系统监控提示选择工具时要考虑团队技术栈避免引入过多新技术增加维护成本。我曾在一个项目中同时使用Flume和Logstash导致运维复杂度陡增。2. 基于Flume的日志采集实战Flume作为Apache顶级项目特别适合高吞吐的日志数据采集场景。下面以最常见的日志文件采集为例分享具体实施步骤。2.1 环境准备与安装首先需要确保Java环境就绪# 检查Java版本 java -version # 应为1.8或以上下载并安装Flumewget https://archive.apache.org/dist/flume/1.9.0/apache-flume-1.9.0-bin.tar.gz tar -zxvf apache-flume-1.9.0-bin.tar.gz cd apache-flume-1.9.0-bin2.2 基础配置示例创建配置文件example.conf# 定义agent的sources/sinks/channels agent.sources r1 agent.sinks k1 agent.channels c1 # 配置source agent.sources.r1.type exec agent.sources.r1.command tail -F /var/log/application.log agent.sources.r1.channels c1 # 配置channel agent.channels.c1.type memory agent.channels.c1.capacity 1000 agent.channels.c1.transactionCapacity 100 # 配置sink agent.sinks.k1.type logger agent.sinks.k1.channel c1启动Flume agentbin/flume-ng agent --conf conf --conf-file example.conf --name agent -Dflume.root.loggerINFO,console2.3 生产环境优化建议在实际生产环境中需要特别注意以下几点Channel选择内存channel性能好但可能丢数据文件channel可靠但IO开销大建议关键数据使用文件channelSink配置HDFS Sink需配置滚动策略Kafka Sink要注意批次大小多个Sink可配置sink组实现负载均衡监控配置启用JMX监控集成Prometheus监控设置关键指标告警我曾遇到一个案例内存channel容量设置过小导致数据丢失后来改用文件channel并适当调整事务容量解决了问题。3. Python数据采集方案对于需要高度定制化的采集场景Python凭借丰富的库成为首选。以下是几种典型实现方式。3.1 基础爬虫实现使用requestsBeautifulSoup的基础爬虫import requests from bs4 import BeautifulSoup url https://example.com/data headers {User-Agent: Mozilla/5.0} response requests.get(url, headersheaders) soup BeautifulSoup(response.text, html.parser) # 提取数据示例 data_items [] for item in soup.select(.data-row): data { title: item.select_one(.title).text, value: float(item.select_one(.value).text) } data_items.append(data)3.2 高级爬虫技巧反爬应对策略使用随机User-Agent设置合理的请求间隔维护IP代理池处理验证码可接入打码平台数据存储优化增量采集记录最后采集位置异常重试实现指数退避重试分布式采集使用Scrapy-Redis性能优化异步请求aiohttp多线程/协程请求批处理3.3 数据处理管道采集后的数据通常需要清洗和转换import pandas as pd # 数据清洗 def clean_data(df): # 处理缺失值 df df.fillna(methodffill) # 去除重复 df df.drop_duplicates() # 类型转换 df[value] pd.to_numeric(df[value], errorscoerce) return df # 数据增强 def enrich_data(df): df[timestamp] pd.to_datetime(now) df[value_category] pd.cut(df[value], bins[0,50,100,np.inf]) return df4. 数据采集系统监控与维护无论使用哪种采集方案完善的监控体系都必不可少。4.1 关键监控指标指标类别具体指标告警阈值监控工具采集量记录数/秒低于历史均值30%Prometheus数据质量空值率5%自定义脚本系统资源CPU使用率70%持续5分钟Grafana延迟端到端延迟1分钟ELK错误率采集错误数10/分钟Sentry4.2 Prometheus监控Flume配置在Flume配置中启用JMXagent.sinks.k1.type org.apache.flume.sink.PrometheusSink agent.sinks.k1.host 0.0.0.0 agent.sinks.k1.port 4141Prometheus配置scrape_configs: - job_name: flume static_configs: - targets: [flume-host:4141]4.3 常见问题排查数据丢失问题检查channel容量验证事务配置检查网络稳定性性能瓶颈使用jstack分析线程状态检查GC日志监控IO等待时间数据重复检查source的重复提交验证channel的事务隔离检查sink的重试机制在一个电商项目中我们曾遇到凌晨数据积压问题最终发现是HDFS Sink的滚动策略与NameNode维护窗口冲突导致调整后问题解决。5. 数据采集最佳实践结合多年项目经验总结以下实践要点设计原则至少保证at-least-once语义关键数据要有端到端校验采集与处理解耦容错机制实现死信队列设计幂等写入建立数据补采流程元数据管理记录数据来源维护数据血缘标注采集时间性能调优批量处理替代单条处理合理设置缓冲区并行化处理流程我曾参与的一个金融项目通过实现数据采集的指纹校验机制成功将数据不一致率从0.1%降到0.001%以下。具体做法是对每批数据生成哈希指纹在关键节点进行校验。